use ReentrantReadWriteLock on TrackImpl and all subclasses to increase concurrency while maintaining consistency

This commit is contained in:
Axel Uhl committed 2012-06-28 14:00:39 +02:00
1 parent 71d48e92a3
commit 64f3d4c2dd
15 files changed
+702 -337

No files matched your search

@@ -28,8 +28,13 @@ public class StoredTrackBasedTestWithTrackedRace extends StoredTrackBasedTest {
private void copyTracks(Map<Competitor, DynamicGPSFixTrack<Competitor, GPSFixMoving>> tracks) {
for (Map.Entry<Competitor, DynamicGPSFixTrack<Competitor, GPSFixMoving>> e : tracks.entrySet()) {
DynamicGPSFixTrack<Competitor, GPSFixMoving> track = getTrackedRace().getTrack(e.getKey());
for (GPSFixMoving fix : e.getValue().getRawFixes()) {
track.addGPSFix(fix);
e.getValue().lockForRead();
try {
for (GPSFixMoving fix : e.getValue().getRawFixes()) {
track.addGPSFix(fix);
}
} finally {
e.getValue().unlockAfterRead();
}
List<MarkPassing> markPassings = new ArrayList<MarkPassing>();
// add a mark passing for the start gate at the very beginning to make sure everyone is on a valid leg
@@ -64,8 +64,13 @@ public class TrackSmootheningTest extends StoredTrackBasedTest {
}
protected void assertOutlierInTrack(DynamicGPSFixTrack<Competitor, GPSFixMoving> track) {
GPSFixMoving outlier = getAnyOutlier(track.getRawFixes());
assertNotNull(outlier); // assert that we found an outlier
track.lockForRead();
try {
GPSFixMoving outlier = getAnyOutlier(track.getRawFixes());
assertNotNull(outlier); // assert that we found an outlier
} finally {
track.unlockAfterRead();
}
}
protected void assertNoOutlierInSmoothenedTrack(DynamicGPSFixTrack<Competitor, GPSFixMoving> track) {
@@ -70,11 +70,17 @@ public class Simulator {
new Thread("Wind simulator for wind source "+windSourceAndTrack.getKey()+" for tracked race "+trackedRace.getRace().getName()) {
@Override
public void run() {
for (Wind wind : windSourceAndTrack.getValue().getRawFixes()) {
if (stopped) {
break;
final WindTrack windTrack = windSourceAndTrack.getValue();
windTrack.lockForRead();
try {
for (Wind wind : windTrack.getRawFixes()) {
if (stopped) {
break;
}
trackedRace.recordWind(delayWind(wind), windSourceAndTrack.getKey());
}
trackedRace.recordWind(delayWind(wind), windSourceAndTrack.getKey());
} finally {
windTrack.unlockAfterRead();
}
}
}.start();
@@ -48,7 +48,7 @@ public class DouglasPeucker<ItemType, FixType extends GPSFix> {
}
private Pair<GPSFix, Distance> getFixWithGreatestCrossTrackErrorInInterval(TimePoint from, TimePoint to) {
Distance maxDistance = Distance.NULL;
Distance maxDistance = Distance.NULL;
FixType firstFixAtOrAfter = track.getFirstFixAtOrAfter(from);
Pair<GPSFix, Distance> result = null;
if (firstFixAtOrAfter != null) {
@@ -56,7 +56,8 @@ public class DouglasPeucker<ItemType, FixType extends GPSFix> {
FixType toFix = track.getLastFixAtOrBefore(to);
if (toFix != null) {
final Bearing bearing = fromPosition.getBearingGreatCircle(toFix.getPosition());
synchronized (track) {
track.lockForRead();
try {
Iterator<FixType> fixIter = track.getFixesIterator(from, /* inclusive */false);
if (executor != null) {
result = getFixWithGreatestCrossTrackErrorUsingExecutor(to, maxDistance, fromPosition, bearing,
@@ -77,6 +78,8 @@ public class DouglasPeucker<ItemType, FixType extends GPSFix> {
}
result = new Pair<GPSFix, Distance>(fixFurthestAway, maxDistance);
}
} finally {
track.unlockAfterRead();
}
}
}
@@ -2,6 +2,7 @@ package com.sap.sailing.domain.tracking;
import java.io.Serializable;
import java.util.Iterator;
import java.util.concurrent.locks.ReadWriteLock;
import com.sap.sailing.domain.base.Timed;
import com.sap.sailing.domain.common.TimePoint;
@@ -12,24 +13,46 @@ import com.sap.sailing.domain.common.TimePoint;
* understanding of how to eliminate outliers. For example, if a track implementation knows it's tracking boats, it may
* consider fixes that the boat cannot possibly have reached due to its speed and direction change limitations as
* outliers. The set of fixes with outliers filtered out can be obtained using {@link #getFixes} whereas
* {@link #getRawFixes()} returns the unfilteres, raw fixes. If an implementation has no idea what an outlier is,
* both methods will return the same fix sequence.
* {@link #getRawFixes()} returns the unfiltered, raw fixes. If an implementation has no idea what an outlier is,
* both methods will return the same fix sequence.<p>
*
* With tracks, concurrency is an important issue. Threads may want to modify a track while other threads may want to
* read from it. Several methods such as {@link #getLastFixAtOrBefore(TimePoint)} return a single fix and can manage
* concurrency internally. However, those methods returning a collection of fixes, such as {@link #getFixes()} or an
* iterator over a collection of fixes, such as {@link #getFixesIterator(TimePoint, boolean)}, need special treatment.
* Until we internalize such iterations (see bug 824, http://bugzilla.sapsailing.com/bugzilla/show_bug.cgi?id=824),
* callers need to manage a read lock which is part of a {@link ReadWriteLock} managed by this track. Callers do so
* by calling {@link #lockForRead} and {@link #unlockAfterRead}.
*
* @author Axel Uhl (d043530)
*/
public interface Track<FixType extends Timed> extends Serializable {
/**
* Callers must synchronize on this object before iterating the result if they have to expect concurrent
* modifications.
* Locks this track for reading by the calling thread. If the thread already holds the lock for this track,
* the hold count will be incremented. Make sure to call {@link #unlockAfterRead()} in a <code>finally</code>
* block to release the lock under all possible circumstances. Failure to do so will inevitably lead to
* deadlocks!
*/
void lockForRead();
/**
* Decrements the hold count for this track's read lock for the calling thread. If it goes to zero, the lock will be
* released and other readers or a writer can obtain the lock. Make sure to call this method in a
* <code>finally</code> block for each {@link #lockForRead()} invocation.
*/
void unlockAfterRead();
/**
* Callers must have called {@link #lockForRead()} before calling this method. This will be checked, and an exception
* will be thrown in case the caller has failed to do so.
*
* @return the raw fixes as recorded by this track; in particular, no smoothening or dampening of any kind is
* applied to the fixes returned by this method.
* @return the smoothened fixes
*/
Iterable<FixType> getFixes();
/**
* Callers must synchronize on this object before iterating the result if they have to expect concurrent
* modifications.
* Callers must have called {@link #lockForRead()} before calling this method. This will be checked, and an exception
* will be thrown in case the caller has failed to do so.
*/
Iterable<FixType> getRawFixes();
@@ -66,8 +89,8 @@ public interface Track<FixType extends Timed> extends Serializable {
* <code>inclusive</code> is <code>true</code>). The fixes returned by the iterator are the smoothened fixes (see
* also {@link #getFixes()}, without any smoothening or dampening applied.
*
* Callers must synchronize on this object before iterating the result if they have to expect concurrent
* modifications.
* Callers must have called {@link #lockForRead()} before calling this method. This will be checked, and an exception
* will be thrown in case the caller has failed to do so.
*/
Iterator<FixType> getFixesIterator(TimePoint startingAt, boolean inclusive);
@@ -76,8 +99,8 @@ public interface Track<FixType extends Timed> extends Serializable {
* <code>inclusive</code> is <code>true</code>). The fixes returned by the iterator are the raw fixes (see also
* {@link #getRawFixes()}, without any smoothening or dampening applied.
*
* Callers must synchronize on this object before iterating the result if they have to expect concurrent
* modifications.
* Callers must have called {@link #lockForRead()} before calling this method. This will be checked, and an exception
* will be thrown in case the caller has failed to do so.
*/
Iterator<FixType> getRawFixesIterator(TimePoint startingAt, boolean inclusive);
@@ -40,6 +40,7 @@ public class CourseBasedWindTrackImpl extends WindTrackImpl {
@Override
protected NavigableSet<Wind> getInternalRawFixes() {
assertReadLock();
NavigableSet<Wind> result;
if (trackedRace.raceIsKnownToStartUpwind()) {
TimePoint startTime = trackedRace.getStartOfRace();
@@ -54,51 +54,66 @@ public class DynamicGPSFixMovingTrackImpl<ItemType> extends DynamicTrackImpl<Ite
@Override
protected SpeedWithBearingWithConfidence<TimePoint> getEstimatedSpeed(TimePoint at,
NavigableSet<GPSFixMoving> fixesToUseForSpeedEstimation, Weigher<TimePoint> weigher) {
List<GPSFixMoving> relevantFixes = getFixesRelevantForSpeedEstimation(at, fixesToUseForSpeedEstimation);
List<SpeedWithConfidence<TimePoint>> speeds = new ArrayList<SpeedWithConfidence<TimePoint>>();
BearingWithConfidenceCluster<TimePoint> bearingCluster = new BearingWithConfidenceCluster<TimePoint>(weigher);
if (!relevantFixes.isEmpty()) {
Iterator<GPSFixMoving> fixIter = relevantFixes.iterator();
GPSFixMoving last = fixIter.next();
SpeedWithConfidenceImpl<TimePoint> speedWithConfidence = new SpeedWithConfidenceImpl<TimePoint>(last.getSpeed(),
/* original confidence */ 0.9, last.getTimePoint());
speeds.add(speedWithConfidence);
bearingCluster.add(new BearingWithConfidenceImpl<TimePoint>(last.getSpeed().getBearing(), /* confidence */ 0.9, last.getTimePoint()));
while (fixIter.hasNext()) {
// add to average the position and time difference
GPSFixMoving next = fixIter.next();
SpeedWithConfidenceImpl<TimePoint> measuredSpeedWithConfidence = new SpeedWithConfidenceImpl<TimePoint>(
last.getPosition().getDistance(next.getPosition())
.inTime(next.getTimePoint().asMillis() - last.getTimePoint().asMillis()),
/* original confidence */0.9, new MillisecondsTimePoint((last.getTimePoint().asMillis()+next.getTimePoint().asMillis())/2));
speeds.add(measuredSpeedWithConfidence);
bearingCluster.add(new BearingWithConfidenceImpl<TimePoint>(last.getPosition().getBearingGreatCircle(next.getPosition()),
/* confidence */ weigher.getConfidence(last.getTimePoint(), next.getTimePoint()),
new MillisecondsTimePoint((last.getTimePoint().asMillis()+next.getTimePoint().asMillis())/2)));
// add to average the speed and bearing provided by the GPSFixMoving
SpeedWithConfidenceImpl<TimePoint> computedSpeedWithConfidence = new SpeedWithConfidenceImpl<TimePoint>(
next.getSpeed(), /* original confidence */0.9, next.getTimePoint());
speeds.add(computedSpeedWithConfidence);
bearingCluster.add(new BearingWithConfidenceImpl<TimePoint>(next.getSpeed().getBearing(), /* confidence */ 0.9, next.getTimePoint()));
last = next;
lockForRead();
try {
List<GPSFixMoving> relevantFixes = getFixesRelevantForSpeedEstimation(at, fixesToUseForSpeedEstimation);
List<SpeedWithConfidence<TimePoint>> speeds = new ArrayList<SpeedWithConfidence<TimePoint>>();
BearingWithConfidenceCluster<TimePoint> bearingCluster = new BearingWithConfidenceCluster<TimePoint>(
weigher);
if (!relevantFixes.isEmpty()) {
Iterator<GPSFixMoving> fixIter = relevantFixes.iterator();
GPSFixMoving last = fixIter.next();
SpeedWithConfidenceImpl<TimePoint> speedWithConfidence = new SpeedWithConfidenceImpl<TimePoint>(
last.getSpeed(),
/* original confidence */0.9, last.getTimePoint());
speeds.add(speedWithConfidence);
bearingCluster.add(new BearingWithConfidenceImpl<TimePoint>(last.getSpeed().getBearing(), /* confidence */
0.9, last.getTimePoint()));
while (fixIter.hasNext()) {
// add to average the position and time difference
GPSFixMoving next = fixIter.next();
SpeedWithConfidenceImpl<TimePoint> measuredSpeedWithConfidence = new SpeedWithConfidenceImpl<TimePoint>(
last.getPosition().getDistance(next.getPosition())
.inTime(next.getTimePoint().asMillis() - last.getTimePoint().asMillis()),
/* original confidence */0.9, new MillisecondsTimePoint(
(last.getTimePoint().asMillis() + next.getTimePoint().asMillis()) / 2));
speeds.add(measuredSpeedWithConfidence);
bearingCluster.add(new BearingWithConfidenceImpl<TimePoint>(last.getPosition()
.getBearingGreatCircle(next.getPosition()),
/* confidence */weigher.getConfidence(last.getTimePoint(), next.getTimePoint()),
new MillisecondsTimePoint(
(last.getTimePoint().asMillis() + next.getTimePoint().asMillis()) / 2)));
// add to average the speed and bearing provided by the GPSFixMoving
SpeedWithConfidenceImpl<TimePoint> computedSpeedWithConfidence = new SpeedWithConfidenceImpl<TimePoint>(
next.getSpeed(), /* original confidence */0.9, next.getTimePoint());
speeds.add(computedSpeedWithConfidence);
bearingCluster.add(new BearingWithConfidenceImpl<TimePoint>(next.getSpeed().getBearing(), /* confidence */
0.9, next.getTimePoint()));
last = next;
}
}
ConfidenceBasedAverager<Double, Speed, TimePoint> speedAverager = ConfidenceFactory.INSTANCE
.createAverager(weigher);
HasConfidence<Double, Speed, TimePoint> speedWithConfidence = speedAverager.getAverage(speeds, at);
BearingWithConfidence<TimePoint> bearingAverage = bearingCluster.getAverage(at);
Bearing bearing = bearingAverage == null ? null : bearingAverage.getObject();
SpeedWithBearing avgSpeed = (speedWithConfidence == null || bearing == null) ? null
: new KnotSpeedWithBearingImpl(speedWithConfidence.getObject().getKnots(), bearing);
SpeedWithBearingWithConfidence<TimePoint> result = speedWithConfidence == null || bearingAverage == null ? null
: new SpeedWithBearingWithConfidenceImpl<TimePoint>(
avgSpeed,
/* confidence */((speedWithConfidence == null ? 0.0 : speedWithConfidence.getConfidence()) + (bearingAverage == null ? 0.0
: bearingAverage.getConfidence())) / 2., at);
return result;
} finally {
unlockAfterRead();
}
ConfidenceBasedAverager<Double, Speed, TimePoint> speedAverager = ConfidenceFactory.INSTANCE.createAverager(weigher);
HasConfidence<Double, Speed, TimePoint> speedWithConfidence = speedAverager.getAverage(speeds, at);
BearingWithConfidence<TimePoint> bearingAverage = bearingCluster.getAverage(at);
Bearing bearing = bearingAverage == null ? null : bearingAverage.getObject();
SpeedWithBearing avgSpeed = (speedWithConfidence == null || bearing == null) ? null :
new KnotSpeedWithBearingImpl(speedWithConfidence.getObject().getKnots(), bearing);
SpeedWithBearingWithConfidence<TimePoint> result = speedWithConfidence == null || bearingAverage == null ? null :
new SpeedWithBearingWithConfidenceImpl<TimePoint>(avgSpeed,
/* confidence */ ((speedWithConfidence == null ? 0.0 : speedWithConfidence.getConfidence()) +
(bearingAverage == null ? 0.0 : bearingAverage.getConfidence()))/2., at);
return result;
}
private List<GPSFixMoving> getFixesRelevantForSpeedEstimation(TimePoint at,
NavigableSet<GPSFixMoving> fixesToUseForSpeedEstimation) {
assertReadLock();
// TODO factor out the obtaining of relevant fixes which should be the same in super.getEstimatedSpeed(at)
DummyGPSFixMoving atTimed = new DummyGPSFixMoving(at);
List<GPSFixMoving> relevantFixes = new LinkedList<GPSFixMoving>();
@@ -189,13 +204,19 @@ public class DynamicGPSFixMovingTrackImpl<ItemType> extends DynamicTrackImpl<Ite
@Override
protected Speed getSpeed(GPSFixMoving fix, Position lastPos, TimePoint timePointOfLastPos) {
Speed fixSpeed = fix.getSpeed();
Speed calculatedSpeed = super.getSpeed(fix, lastPos, timePointOfLastPos);
Speed averaged = averageSpeed(fixSpeed, calculatedSpeed);
return averaged;
lockForRead();
try {
Speed fixSpeed = fix.getSpeed();
Speed calculatedSpeed = super.getSpeed(fix, lastPos, timePointOfLastPos);
Speed averaged = averageSpeed(fixSpeed, calculatedSpeed);
return averaged;
} finally {
unlockAfterRead();
}
}
private Speed averageSpeed(Speed... speeds) {
assertReadLock();
double sumInKMH = 0;
int count = 0;
for (Speed speed : speeds) {
@@ -214,6 +235,7 @@ public class DynamicGPSFixMovingTrackImpl<ItemType> extends DynamicTrackImpl<Ite
*/
@Override
protected boolean isValid(PartialNavigableSetView<GPSFixMoving> filteredView, GPSFixMoving e) {
assertReadLock();
boolean result;
if (e.isValidityCached()) {
result = e.isValid();
@@ -19,9 +19,12 @@ public class DynamicTrackImpl<ItemType, FixType extends GPSFix> extends
@Override
public void addGPSFix(FixType gpsFix) {
synchronized (this) {
lockForWrite();
try {
getInternalRawFixes().add(gpsFix);
invalidateValidityCaches(gpsFix);
} finally {
unlockAfterWrite();
}
Iterable<GPSTrackListener<ItemType, FixType>> listeners = getListeners();
synchronized (listeners) {
@@ -466,8 +466,13 @@ public class DynamicTrackedRaceImpl extends TrackedRaceImpl implements
WindTrack result = super.createWindTrack(windSource, delayForWindEstimationCacheInvalidation);
if (windSource.getType().canBeStored()) {
// replicate all wind fixed that may have been loaded by the wind store
for (Wind wind : result.getRawFixes()) {
notifyListeners(wind, windSource);
result.lockForRead();
try {
for (Wind wind : result.getRawFixes()) {
notifyListeners(wind, windSource);
}
} finally {
result.unlockAfterRead();
}
}
return result;
@@ -137,81 +137,116 @@ public class GPSFixTrackImpl<ItemType, FixType extends GPSFix> extends TrackImpl
}
private Pair<FixType, FixType> getFixesForPositionEstimation(TimePoint timePoint, boolean inclusive) {
FixType lastFixBefore = inclusive ? getLastFixAtOrBefore(timePoint) : getLastFixBefore(timePoint);
FixType firstFixAfter = inclusive ? getFirstFixAtOrAfter(timePoint) : getFirstFixAfter(timePoint);
return new Pair<FixType, FixType>(lastFixBefore, firstFixAfter);
lockForRead();
try {
FixType lastFixBefore = inclusive ? getLastFixAtOrBefore(timePoint) : getLastFixBefore(timePoint);
FixType firstFixAfter = inclusive ? getFirstFixAtOrAfter(timePoint) : getFirstFixAfter(timePoint);
return new Pair<FixType, FixType>(lastFixBefore, firstFixAfter);
} finally {
unlockAfterRead();
}
}
@Override
public Position getEstimatedPosition(TimePoint timePoint, boolean extrapolate) {
Pair<FixType, FixType> fixesForPositionEstimation = getFixesForPositionEstimation(timePoint, /* inclusive */ true);
return getEstimatedPosition(timePoint, extrapolate, fixesForPositionEstimation.getA(), fixesForPositionEstimation.getB());
lockForRead();
try {
Pair<FixType, FixType> fixesForPositionEstimation = getFixesForPositionEstimation(timePoint, /* inclusive */ true);
return getEstimatedPosition(timePoint, extrapolate, fixesForPositionEstimation.getA(), fixesForPositionEstimation.getB());
} finally {
unlockAfterRead();
}
}
@Override
public Pair<TimePoint, TimePoint> getEstimatedPositionTimePeriodAffectedBy(GPSFix fix) {
Pair<FixType, FixType> fixesForPositionEstimation = getFixesForPositionEstimation(fix.getTimePoint(), /* inclusive */ false);
return new Pair<TimePoint, TimePoint>(fixesForPositionEstimation.getA() == null ? null : fixesForPositionEstimation.getA().getTimePoint(),
fixesForPositionEstimation.getB() == null ? null : fixesForPositionEstimation.getB().getTimePoint());
lockForRead();
try {
Pair<FixType, FixType> fixesForPositionEstimation = getFixesForPositionEstimation(fix.getTimePoint(), /* inclusive */ false);
return new Pair<TimePoint, TimePoint>(fixesForPositionEstimation.getA() == null ? null : fixesForPositionEstimation.getA().getTimePoint(),
fixesForPositionEstimation.getB() == null ? null : fixesForPositionEstimation.getB().getTimePoint());
} finally {
unlockAfterRead();
}
}
@Override
public Position getEstimatedRawPosition(TimePoint timePoint, boolean extrapolate) {
FixType lastFixAtOrBefore = getLastRawFixAtOrBefore(timePoint);
FixType firstFixAtOrAfter = getFirstRawFixAtOrAfter(timePoint);
return getEstimatedPosition(timePoint, extrapolate, lastFixAtOrBefore, firstFixAtOrAfter);
lockForRead();
try {
FixType lastFixAtOrBefore = getLastRawFixAtOrBefore(timePoint);
FixType firstFixAtOrAfter = getFirstRawFixAtOrAfter(timePoint);
return getEstimatedPosition(timePoint, extrapolate, lastFixAtOrBefore, firstFixAtOrAfter);
} finally {
unlockAfterRead();
}
}
private Position getEstimatedPosition(TimePoint timePoint, boolean extrapolate, FixType lastFixAtOrBefore,
FixType firstFixAtOrAfter) {
// TODO bug #346: compute a confidence value for the position returned based on time difference between fix(es) and timePoint; consider using Taylor approximation of more fixes around timePoint to predict and weigh position
if (lastFixAtOrBefore != null && lastFixAtOrBefore == firstFixAtOrAfter) {
return lastFixAtOrBefore.getPosition(); // exact match; how unlikely is that?
} else {
if (lastFixAtOrBefore == null && firstFixAtOrAfter != null) {
return firstFixAtOrAfter.getPosition(); // asking for time point before first fix: return first fix's position
}
if (firstFixAtOrAfter == null && !extrapolate) {
return lastFixAtOrBefore == null ? null : lastFixAtOrBefore.getPosition();
lockForRead();
try {
// TODO bug #346: compute a confidence value for the position returned based on time difference between fix(es) and timePoint; consider using Taylor approximation of more fixes around timePoint to predict and weigh position
if (lastFixAtOrBefore != null && lastFixAtOrBefore == firstFixAtOrAfter) {
return lastFixAtOrBefore.getPosition(); // exact match; how unlikely is that?
} else {
SpeedWithBearing estimatedSpeed = estimateSpeed(lastFixAtOrBefore, firstFixAtOrAfter);
if (estimatedSpeed == null) {
return null;
if (lastFixAtOrBefore == null && firstFixAtOrAfter != null) {
return firstFixAtOrAfter.getPosition(); // asking for time point before first fix: return first fix's position
}
if (firstFixAtOrAfter == null && !extrapolate) {
return lastFixAtOrBefore == null ? null : lastFixAtOrBefore.getPosition();
} else {
if (lastFixAtOrBefore != null) {
Distance distance = estimatedSpeed.travel(lastFixAtOrBefore.getTimePoint(), timePoint);
Position result = lastFixAtOrBefore.getPosition().translateGreatCircle(
estimatedSpeed.getBearing(), distance);
return result;
SpeedWithBearing estimatedSpeed = estimateSpeed(lastFixAtOrBefore, firstFixAtOrAfter);
if (estimatedSpeed == null) {
return null;
} else {
// firstFixAtOrAfter can't be null because otherwise no speed could have been estimated
return firstFixAtOrAfter.getPosition();
if (lastFixAtOrBefore != null) {
Distance distance = estimatedSpeed.travel(lastFixAtOrBefore.getTimePoint(), timePoint);
Position result = lastFixAtOrBefore.getPosition().translateGreatCircle(
estimatedSpeed.getBearing(), distance);
return result;
} else {
// firstFixAtOrAfter can't be null because otherwise no speed could have been estimated
return firstFixAtOrAfter.getPosition();
}
}
}
}
} finally {
unlockAfterRead();
}
}
@Override
public synchronized Speed getMaximumSpeedOverGround(TimePoint from, TimePoint to) {
// fetch all fixes on this leg so far and determine their maximum speed
Iterator<FixType> iter = getFixesIterator(from, /* inclusive */ true);
Speed max = Speed.NULL;
if (iter.hasNext()) {
Position lastPos = getEstimatedPosition(from, false);
while (iter.hasNext()) {
FixType fix = iter.next();
Speed fixSpeed = getSpeed(fix, lastPos, from);
if (fixSpeed.compareTo(max) > 0) {
max = fixSpeed;
public Speed getMaximumSpeedOverGround(TimePoint from, TimePoint to) {
lockForRead();
try {
// fetch all fixes on this leg so far and determine their maximum speed
Iterator<FixType> iter = getFixesIterator(from, /* inclusive */ true);
Speed max = Speed.NULL;
if (iter.hasNext()) {
Position lastPos = getEstimatedPosition(from, false);
while (iter.hasNext()) {
FixType fix = iter.next();
Speed fixSpeed = getSpeed(fix, lastPos, from);
if (fixSpeed.compareTo(max) > 0) {
max = fixSpeed;
}
}
}
return max;
} finally {
unlockAfterRead();
}
return max;
}
protected Speed getSpeed(FixType fix, Position lastPos, TimePoint timePointOfLastPos) {
return lastPos.getDistance(fix.getPosition()).inTime(fix.getTimePoint().asMillis()-timePointOfLastPos.asMillis());
lockForRead();
try {
return lastPos.getDistance(fix.getPosition()).inTime(fix.getTimePoint().asMillis()-timePointOfLastPos.asMillis());
} finally {
unlockAfterRead();
}
}
private SpeedWithBearing estimateSpeed(FixType fix1, FixType fix2) {
@@ -254,51 +289,61 @@ public class GPSFixTrackImpl<ItemType, FixType extends GPSFix> extends TrackImpl
@Override
public Distance getDistanceTraveled(TimePoint from, TimePoint to) {
double distanceInNauticalMiles = 0;
if (from.compareTo(to) < 0) {
Position fromPos = getEstimatedPosition(from, false);
if (fromPos == null) {
lockForRead();
try {
double distanceInNauticalMiles = 0;
if (from.compareTo(to) < 0) {
Position fromPos = getEstimatedPosition(from, false);
if (fromPos == null) {
return Distance.NULL;
}
synchronized (this) {
NavigableSet<GPSFix> subset = getGPSFixes().subSet(new DummyGPSFix(from),
/* fromInclusive */false, new DummyGPSFix(to),
/* toInclusive */false);
for (GPSFix fix : subset) {
double distanceBetweenAdjacentFixesInNauticalMiles = fromPos.getDistance(fix.getPosition()).getNauticalMiles();
distanceInNauticalMiles += distanceBetweenAdjacentFixesInNauticalMiles;
fromPos = fix.getPosition();
}
}
Position toPos = getEstimatedPosition(to, false);
distanceInNauticalMiles += fromPos.getDistance(toPos).getNauticalMiles();
return new NauticalMileDistance(distanceInNauticalMiles);
} else {
return Distance.NULL;
}
synchronized (this) {
NavigableSet<GPSFix> subset = getGPSFixes().subSet(new DummyGPSFix(from),
/* fromInclusive */false, new DummyGPSFix(to),
/* toInclusive */false);
for (GPSFix fix : subset) {
double distanceBetweenAdjacentFixesInNauticalMiles = fromPos.getDistance(fix.getPosition()).getNauticalMiles();
distanceInNauticalMiles += distanceBetweenAdjacentFixesInNauticalMiles;
fromPos = fix.getPosition();
}
}
Position toPos = getEstimatedPosition(to, false);
distanceInNauticalMiles += fromPos.getDistance(toPos).getNauticalMiles();
return new NauticalMileDistance(distanceInNauticalMiles);
} else {
return Distance.NULL;
} finally {
unlockAfterRead();
}
}
@Override
public Distance getRawDistanceTraveled(TimePoint from, TimePoint to) {
double distanceInNauticalMiles = 0;
if (from.compareTo(to) < 0) {
Position fromPos = getEstimatedRawPosition(from, false);
if (fromPos == null) {
lockForRead();
try {
double distanceInNauticalMiles = 0;
if (from.compareTo(to) < 0) {
Position fromPos = getEstimatedRawPosition(from, false);
if (fromPos == null) {
return Distance.NULL;
}
@SuppressWarnings("unchecked")
NavigableSet<GPSFix> subset = (NavigableSet<GPSFix>) getInternalRawFixes().subSet((FixType) new DummyGPSFix(from),
/* fromInclusive */false, (FixType) new DummyGPSFix(to),
/* toInclusive */false);
for (GPSFix fix : subset) {
distanceInNauticalMiles += fromPos.getDistance(fix.getPosition()).getNauticalMiles();
fromPos = fix.getPosition();
}
Position toPos = getEstimatedRawPosition(to, false);
distanceInNauticalMiles += fromPos.getDistance(toPos).getNauticalMiles();
return new NauticalMileDistance(distanceInNauticalMiles);
} else {
return Distance.NULL;
}
@SuppressWarnings("unchecked")
NavigableSet<GPSFix> subset = (NavigableSet<GPSFix>) getInternalRawFixes().subSet((FixType) new DummyGPSFix(from),
/* fromInclusive */false, (FixType) new DummyGPSFix(to),
/* toInclusive */false);
for (GPSFix fix : subset) {
distanceInNauticalMiles += fromPos.getDistance(fix.getPosition()).getNauticalMiles();
fromPos = fix.getPosition();
}
Position toPos = getEstimatedRawPosition(to, false);
distanceInNauticalMiles += fromPos.getDistance(toPos).getNauticalMiles();
return new NauticalMileDistance(distanceInNauticalMiles);
} else {
return Distance.NULL;
} finally {
unlockAfterRead();
}
}
@@ -311,29 +356,44 @@ public class GPSFixTrackImpl<ItemType, FixType extends GPSFix> extends TrackImpl
* the speeds/bearings determined by distance/time difference of the fixes themselves.
*/
@Override
public synchronized SpeedWithBearing getEstimatedSpeed(TimePoint at) {
SpeedWithBearingWithConfidence<TimePoint> estimatedSpeed = getEstimatedSpeed(at, getInternalFixes(),
ConfidenceFactory.INSTANCE.createExponentialTimeDifferenceWeigher(
// use a minimum confidence to avoid the bearing to flip to 270deg in case all is zero
getMillisecondsOverWhichToAverageSpeed()));
return estimatedSpeed == null ? null : estimatedSpeed.getObject();
public SpeedWithBearing getEstimatedSpeed(TimePoint at) {
lockForRead();
try {
SpeedWithBearingWithConfidence<TimePoint> estimatedSpeed = getEstimatedSpeed(at, getInternalFixes(),
ConfidenceFactory.INSTANCE.createExponentialTimeDifferenceWeigher(
// use a minimum confidence to avoid the bearing to flip to 270deg in case all is zero
getMillisecondsOverWhichToAverageSpeed()));
return estimatedSpeed == null ? null : estimatedSpeed.getObject();
} finally {
unlockAfterRead();
}
}
@Override
public synchronized SpeedWithBearing getRawEstimatedSpeed(TimePoint at) {
public SpeedWithBearing getRawEstimatedSpeed(TimePoint at) {
return getEstimatedSpeed(at, getRawFixes(), ConfidenceFactory.INSTANCE.createExponentialTimeDifferenceWeigher(
// use a minimum confidence to avoid the bearing to flip to 270deg in case all is zero
// use a minimum confidence to avoid the bearing to flip to 270deg in case all is zero
getMillisecondsOverWhichToAverageSpeed())).getObject();
}
@Override
public SpeedWithBearingWithConfidence<TimePoint> getEstimatedSpeed(TimePoint at, Weigher<TimePoint> weigher) {
return getEstimatedSpeed(at, getInternalFixes(), weigher);
lockForRead();
try {
return getEstimatedSpeed(at, getInternalFixes(), weigher);
} finally {
unlockAfterRead();
}
}
@Override
public SpeedWithBearingWithConfidence<TimePoint> getRawEstimatedSpeed(TimePoint at, Weigher<TimePoint> weigher) {
return getEstimatedSpeed(at, getRawFixes(), weigher);
lockForRead();
try {
return getEstimatedSpeed(at, getRawFixes(), weigher);
} finally {
unlockAfterRead();
}
}
/**
@@ -351,62 +411,72 @@ public class GPSFixTrackImpl<ItemType, FixType extends GPSFix> extends TrackImpl
*/
protected SpeedWithBearingWithConfidence<TimePoint> getEstimatedSpeed(TimePoint at,
NavigableSet<FixType> fixesToUseForSpeedEstimation, Weigher<TimePoint> weigher) {
@SuppressWarnings("unchecked")
NavigableSet<GPSFix> gpsFixesToUseForSpeedEstimation = (NavigableSet<GPSFix>) fixesToUseForSpeedEstimation;
List<GPSFix> relevantFixes = getFixesRelevantForSpeedEstimation(at, gpsFixesToUseForSpeedEstimation);
List<SpeedWithConfidence<TimePoint>> speeds = new ArrayList<SpeedWithConfidence<TimePoint>>();
BearingWithConfidenceCluster<TimePoint> bearingCluster = new BearingWithConfidenceCluster<TimePoint>(weigher);
if (!relevantFixes.isEmpty()) {
Iterator<GPSFix> fixIter = relevantFixes.iterator();
GPSFix last = fixIter.next();
while (fixIter.hasNext()) {
// TODO bug #346: consider time difference between next.getTimepoint() and at to compute a confidence
GPSFix next = fixIter.next();
// TODO bug #345: use SpeedWithConfidence to aggregate confidence-tagged speed values
MillisecondsTimePoint relativeTo = new MillisecondsTimePoint((last.getTimePoint().asMillis() + next.getTimePoint().asMillis())/2);
Speed speed = last.getPosition().getDistance(next.getPosition())
.inTime(next.getTimePoint().asMillis() - last.getTimePoint().asMillis());
SpeedWithConfidenceImpl<TimePoint> speedWithConfidence = new SpeedWithConfidenceImpl<TimePoint>(speed, /* original confidence */
0.9, relativeTo);
speeds.add(speedWithConfidence);
bearingCluster.add(new BearingWithConfidenceImpl<TimePoint>(last.getPosition().getBearingGreatCircle(next.getPosition()),
/* confidence */ 0.9, // TODO use number of tracked satellites to determine confidence of single fix
relativeTo));
last = next;
lockForRead();
try {
@SuppressWarnings("unchecked")
NavigableSet<GPSFix> gpsFixesToUseForSpeedEstimation = (NavigableSet<GPSFix>) fixesToUseForSpeedEstimation;
List<GPSFix> relevantFixes = getFixesRelevantForSpeedEstimation(at, gpsFixesToUseForSpeedEstimation);
List<SpeedWithConfidence<TimePoint>> speeds = new ArrayList<SpeedWithConfidence<TimePoint>>();
BearingWithConfidenceCluster<TimePoint> bearingCluster = new BearingWithConfidenceCluster<TimePoint>(weigher);
if (!relevantFixes.isEmpty()) {
Iterator<GPSFix> fixIter = relevantFixes.iterator();
GPSFix last = fixIter.next();
while (fixIter.hasNext()) {
// TODO bug #346: consider time difference between next.getTimepoint() and at to compute a confidence
GPSFix next = fixIter.next();
// TODO bug #345: use SpeedWithConfidence to aggregate confidence-tagged speed values
MillisecondsTimePoint relativeTo = new MillisecondsTimePoint((last.getTimePoint().asMillis() + next.getTimePoint().asMillis())/2);
Speed speed = last.getPosition().getDistance(next.getPosition())
.inTime(next.getTimePoint().asMillis() - last.getTimePoint().asMillis());
SpeedWithConfidenceImpl<TimePoint> speedWithConfidence = new SpeedWithConfidenceImpl<TimePoint>(speed, /* original confidence */
0.9, relativeTo);
speeds.add(speedWithConfidence);
bearingCluster.add(new BearingWithConfidenceImpl<TimePoint>(last.getPosition().getBearingGreatCircle(next.getPosition()),
/* confidence */ 0.9, // TODO use number of tracked satellites to determine confidence of single fix
relativeTo));
last = next;
}
}
ConfidenceBasedAverager<Double, Speed, TimePoint> speedAverager = ConfidenceFactory.INSTANCE.createAverager(weigher);
HasConfidence<Double, Speed, TimePoint> speedWithConfidence = speedAverager.getAverage(speeds, at);
BearingWithConfidence<TimePoint> bearingAverage = bearingCluster.getAverage(at);
Bearing bearing = bearingAverage == null ? null : bearingAverage.getObject();
SpeedWithBearing avgSpeed = (speedWithConfidence == null || bearing == null) ? null :
new KnotSpeedWithBearingImpl(speedWithConfidence.getObject().getKnots(), bearing);
SpeedWithBearingWithConfidence<TimePoint> result = avgSpeed == null ? null :
new SpeedWithBearingWithConfidenceImpl<TimePoint>(avgSpeed, (bearingAverage.getConfidence() + speedWithConfidence.getConfidence())/2., at);
return result;
} finally {
unlockAfterRead();
}
ConfidenceBasedAverager<Double, Speed, TimePoint> speedAverager = ConfidenceFactory.INSTANCE.createAverager(weigher);
HasConfidence<Double, Speed, TimePoint> speedWithConfidence = speedAverager.getAverage(speeds, at);
BearingWithConfidence<TimePoint> bearingAverage = bearingCluster.getAverage(at);
Bearing bearing = bearingAverage == null ? null : bearingAverage.getObject();
SpeedWithBearing avgSpeed = (speedWithConfidence == null || bearing == null) ? null :
new KnotSpeedWithBearingImpl(speedWithConfidence.getObject().getKnots(), bearing);
SpeedWithBearingWithConfidence<TimePoint> result = avgSpeed == null ? null :
new SpeedWithBearingWithConfidenceImpl<TimePoint>(avgSpeed, (bearingAverage.getConfidence() + speedWithConfidence.getConfidence())/2., at);
return result;
}
private List<GPSFix> getFixesRelevantForSpeedEstimation(TimePoint at,
NavigableSet<GPSFix> fixesToUseForSpeedEstimation) {
DummyGPSFix atTimed = new DummyGPSFix(at);
List<GPSFix> relevantFixes = new LinkedList<GPSFix>();
synchronized (this) {
NavigableSet<GPSFix> beforeSet = fixesToUseForSpeedEstimation.headSet(atTimed, /* inclusive */ false);
NavigableSet<GPSFix> afterSet = fixesToUseForSpeedEstimation.tailSet(atTimed, /* inclusive */ true);
for (GPSFix beforeFix : beforeSet.descendingSet()) {
if (at.asMillis() - beforeFix.getTimePoint().asMillis() > getMillisecondsOverWhichToAverage() / 2) {
break;
lockForRead();
try {
DummyGPSFix atTimed = new DummyGPSFix(at);
List<GPSFix> relevantFixes = new LinkedList<GPSFix>();
synchronized (this) {
NavigableSet<GPSFix> beforeSet = fixesToUseForSpeedEstimation.headSet(atTimed, /* inclusive */ false);
NavigableSet<GPSFix> afterSet = fixesToUseForSpeedEstimation.tailSet(atTimed, /* inclusive */ true);
for (GPSFix beforeFix : beforeSet.descendingSet()) {
if (at.asMillis() - beforeFix.getTimePoint().asMillis() > getMillisecondsOverWhichToAverage() / 2) {
break;
}
relevantFixes.add(0, beforeFix);
}
relevantFixes.add(0, beforeFix);
}
for (GPSFix afterFix : afterSet) {
if (afterFix.getTimePoint().asMillis() - at.asMillis() > getMillisecondsOverWhichToAverage() / 2) {
break;
for (GPSFix afterFix : afterSet) {
if (afterFix.getTimePoint().asMillis() - at.asMillis() > getMillisecondsOverWhichToAverage() / 2) {
break;
}
relevantFixes.add(afterFix);
}
relevantFixes.add(afterFix);
}
return relevantFixes;
} finally {
unlockAfterRead();
}
return relevantFixes;
}
protected long getMillisecondsOverWhichToAverage() {
@@ -418,6 +488,7 @@ public class GPSFixTrackImpl<ItemType, FixType extends GPSFix> extends TrackImpl
*/
@Override
protected NavigableSet<FixType> getInternalFixes() {
assertReadLock();
return new PartialNavigableSetView<FixType>(super.getInternalFixes()) {
@Override
protected boolean isValid(FixType e) {
@@ -432,6 +503,7 @@ public class GPSFixTrackImpl<ItemType, FixType extends GPSFix> extends TrackImpl
* adding a fix, only immediately adjacent fix's validity caches need to be invalidated.
*/
protected boolean isValid(PartialNavigableSetView<FixType> filteredView, FixType e) {
assertReadLock();
boolean result;
if (maxSpeedForSmoothening == null) {
result = true;
@@ -464,7 +536,8 @@ public class GPSFixTrackImpl<ItemType, FixType extends GPSFix> extends TrackImpl
* of the fixes whose validity may be affected. If subclasses redefine {@link #isValid(PartialNavigableSetView, GPSFix)},
* they must make sure that this method is redefined accordingly.
*/
protected synchronized void invalidateValidityCaches(FixType gpsFix) {
protected void invalidateValidityCaches(FixType gpsFix) {
assertWriteLock();
gpsFix.invalidateCache();
FixType lower = getInternalRawFixes().lower(gpsFix);
if (lower != null) {
@@ -478,23 +551,30 @@ public class GPSFixTrackImpl<ItemType, FixType extends GPSFix> extends TrackImpl
@Override
public boolean hasDirectionChange(TimePoint at, double minimumDegreeDifference) {
boolean result = false;
TimePoint start = new MillisecondsTimePoint(at.asMillis()-getMillisecondsOverWhichToAverageSpeed());
TimePoint end = new MillisecondsTimePoint(at.asMillis()+getMillisecondsOverWhichToAverageSpeed());
SpeedWithBearing estimatedSpeedAtStart = getEstimatedSpeed(start);
if (estimatedSpeedAtStart != null) {
Bearing bearingAtStart = estimatedSpeedAtStart.getBearing();
TimePoint next = new MillisecondsTimePoint(start.asMillis()+Math.max(1000l, getMillisecondsOverWhichToAverageSpeed()/2));
while (!result && next.compareTo(end) <= 0) {
SpeedWithBearing estimatedSpeedAtNext = getEstimatedSpeed(next);
if (estimatedSpeedAtNext != null) {
Bearing bearingAtEnd = estimatedSpeedAtNext.getBearing();
result = Math.abs(bearingAtStart.getDifferenceTo(bearingAtEnd).getDegrees()) > minimumDegreeDifference;
lockForRead();
try {
boolean result = false;
TimePoint start = new MillisecondsTimePoint(at.asMillis() - getMillisecondsOverWhichToAverageSpeed());
TimePoint end = new MillisecondsTimePoint(at.asMillis() + getMillisecondsOverWhichToAverageSpeed());
SpeedWithBearing estimatedSpeedAtStart = getEstimatedSpeed(start);
if (estimatedSpeedAtStart != null) {
Bearing bearingAtStart = estimatedSpeedAtStart.getBearing();
TimePoint next = new MillisecondsTimePoint(start.asMillis()
+ Math.max(1000l, getMillisecondsOverWhichToAverageSpeed() / 2));
while (!result && next.compareTo(end) <= 0) {
SpeedWithBearing estimatedSpeedAtNext = getEstimatedSpeed(next);
if (estimatedSpeedAtNext != null) {
Bearing bearingAtEnd = estimatedSpeedAtNext.getBearing();
result = Math.abs(bearingAtStart.getDifferenceTo(bearingAtEnd).getDegrees()) > minimumDegreeDifference;
}
next = new MillisecondsTimePoint(next.asMillis()
+ Math.max(1000l, getMillisecondsOverWhichToAverageSpeed() / 2));
}
next = new MillisecondsTimePoint(next.asMillis()+Math.max(1000l, getMillisecondsOverWhichToAverageSpeed()/2));
}
return result;
} finally {
unlockAfterRead();
}
return result;
}
}
@@ -174,8 +174,13 @@ public class TrackBasedEstimationWindTrackImpl extends VirtualWindTrackImpl impl
/**
* Synchronizes serialization on this object to avoid the cache being updated while being written.
*/
private synchronized void writeObject(ObjectOutputStream s) throws IOException {
s.defaultWriteObject();
private void writeObject(ObjectOutputStream s) throws IOException {
lockForRead();
try {
s.defaultWriteObject();
} finally {
unlockAfterRead();
}
}
/**
@@ -214,13 +219,25 @@ public class TrackBasedEstimationWindTrackImpl extends VirtualWindTrackImpl impl
}
protected void cache(TimePoint timePoint, WindWithConfidence<TimePoint> fix) {
synchronized (scheduledRefreshInterval) {
// can't use lockForWrite() here because caching can happen while holding the read lock, and the lock can't be
// upgraded. But lockForRead() and synchronization will do the job because all invalidations lock the write lock,
// and all contains() checks and get() calls use synchronization too.
lockForRead();
try {
if (fix == null) {
getTimePointsWithCachedNullResult().add(timePoint);
timePointsWithCachedNullResultFastContains.add(timePoint);
synchronized (timePointsWithCachedNullResult) {
timePointsWithCachedNullResult.add(timePoint);
}
synchronized (timePointsWithCachedNullResultFastContains) {
timePointsWithCachedNullResultFastContains.add(timePoint);
}
} else {
getCachedFixes().add(fix);
synchronized (cache) {
cache.add(fix);
}
}
} finally {
unlockAfterRead();
}
}
@@ -230,8 +247,9 @@ public class TrackBasedEstimationWindTrackImpl extends VirtualWindTrackImpl impl
* be bundled, and during live mode the incoming requests for a time point close to the time for which new data is received
* will not be massively delayed by having to re-calculate the estimation over and over again.
*/
private synchronized void scheduleCacheRefresh(WindWithConfidence<TimePoint> startOfInvalidation, TimePoint endOfInvalidation) {
synchronized (scheduledRefreshInterval) {
private void scheduleCacheRefresh(WindWithConfidence<TimePoint> startOfInvalidation, TimePoint endOfInvalidation) {
lockForWrite();
try {
if (!scheduledRefreshInterval.isSet()) {
// according to the invariant this implies [1]==null
scheduledRefreshInterval.set(startOfInvalidation, endOfInvalidation);
@@ -241,6 +259,8 @@ public class TrackBasedEstimationWindTrackImpl extends VirtualWindTrackImpl impl
// we can safely extend the interval; the invalidation won't start before we release the lock
scheduledRefreshInterval.extend(startOfInvalidation, endOfInvalidation);
}
} finally {
unlockAfterWrite();
}
}
@@ -250,7 +270,8 @@ public class TrackBasedEstimationWindTrackImpl extends VirtualWindTrackImpl impl
* running.
*/
private void invalidateCache() {
synchronized (scheduledRefreshInterval) {
lockForWrite();
try {
Iterator<WindWithConfidence<TimePoint>> iter = (scheduledRefreshInterval.getStart() == null ? getCachedFixes()
: getCachedFixes().tailSet(scheduledRefreshInterval.getStart(), /* inclusive */true)).iterator();
while (iter.hasNext()) {
@@ -274,6 +295,8 @@ public class TrackBasedEstimationWindTrackImpl extends VirtualWindTrackImpl impl
}
}
scheduledRefreshInterval.clear();
} finally {
unlockAfterWrite();
}
}
@@ -283,7 +306,8 @@ public class TrackBasedEstimationWindTrackImpl extends VirtualWindTrackImpl impl
private void refreshCacheIncrementally() {
Set<WindWithConfidence<TimePoint>> windFixesToRecalculate = new HashSet<WindWithConfidence<TimePoint>>();
Set<TimePoint> cachedNullResultsToRecalculate = new HashSet<TimePoint>();
synchronized (scheduledRefreshInterval) {
lockForWrite();
try {
Iterator<WindWithConfidence<TimePoint>> iter = (scheduledRefreshInterval.getStart() == null ? getCachedFixes()
: getCachedFixes().tailSet(scheduledRefreshInterval.getStart(), /* inclusive */true)).iterator();
Iterator<TimePoint> nullIter = (scheduledRefreshInterval.getStart() == null ? getTimePointsWithCachedNullResult()
@@ -301,6 +325,8 @@ public class TrackBasedEstimationWindTrackImpl extends VirtualWindTrackImpl impl
scheduledRefreshInterval.getEnd() == null) {
cachedNullResultsToRecalculate.add(nextNullResultToRecalculate);
}
} finally {
unlockAfterWrite();
}
Set<TimePoint> nullRemovals = new HashSet<TimePoint>();
Set<TimePoint> nullInsertions = new HashSet<TimePoint>();
@@ -325,7 +351,8 @@ public class TrackBasedEstimationWindTrackImpl extends VirtualWindTrackImpl impl
}
}
// apply the computed cache deltas
synchronized (scheduledRefreshInterval) {
lockForWrite();
try {
for (TimePoint nullRemoval : nullRemovals) {
getTimePointsWithCachedNullResult().remove(nullRemoval);
timePointsWithCachedNullResultFastContains.remove(nullRemoval);
@@ -339,38 +366,44 @@ public class TrackBasedEstimationWindTrackImpl extends VirtualWindTrackImpl impl
for (Map.Entry<TimePoint, WindWithConfidence<TimePoint>> cacheInsertion : cacheInsertions.entrySet()) {
cache(cacheInsertion.getKey(), cacheInsertion.getValue());
}
scheduledRefreshInterval.clear();
} finally {
unlockAfterWrite();
}
scheduledRefreshInterval.clear();
}
private void startSchedulerForCacheRefresh() {
synchronized (scheduledRefreshInterval) {
if (delayForCacheInvalidationInMilliseconds == 0) {
invalidateCache();
} else {
final Timer cacheInvalidationTimer = new Timer("TrackBasedEstimationWindTrackImpl cache invalidation timer for race "
+ getTrackedRace().getRace());
cacheInvalidationTimer.schedule(new TimerTask() {
@Override
public void run() {
// to avoid deadlock with another invalidateCache() and with scheduleCacheInvalidation we need
// to obtain the TrackBasedEstimationWindTrackImpl.this monitor first (see bug 746).
synchronized (scheduledRefreshInterval) {
cacheInvalidationTimer.cancel(); // terminates the timer thread
refreshCacheIncrementally();
}
assertWriteLock();
if (delayForCacheInvalidationInMilliseconds == 0) {
invalidateCache();
} else {
final Timer cacheInvalidationTimer = new Timer(
"TrackBasedEstimationWindTrackImpl cache invalidation timer for race " + getTrackedRace().getRace());
cacheInvalidationTimer.schedule(new TimerTask() {
@Override
public void run() {
// to avoid deadlock with another invalidateCache() and with scheduleCacheInvalidation we need
// to obtain the TrackBasedEstimationWindTrackImpl.this monitor first (see bug 746).
lockForWrite();
try {
cacheInvalidationTimer.cancel(); // terminates the timer thread
refreshCacheIncrementally();
} finally {
unlockAfterWrite();
}
}, delayForCacheInvalidationInMilliseconds);
}
}
}, delayForCacheInvalidationInMilliseconds);
}
}
private void clearCache() {
synchronized (scheduledRefreshInterval) {
getCachedFixes().clear();
lockForWrite();
try {
cache.clear();
timePointsWithCachedNullResult.clear();
timePointsWithCachedNullResultFastContains.clear();
} finally {
unlockAfterWrite();
}
}
@@ -381,12 +414,23 @@ public class TrackBasedEstimationWindTrackImpl extends VirtualWindTrackImpl impl
* it. The result will then be added to the cache.
*/
private WindWithConfidence<TimePoint> getEstimatedWindDirection(Position p, TimePoint timePoint) {
WindWithConfidence<TimePoint> result;
if (nullResultCacheContains(timePoint)) {
result = null;
} else {
WindWithConfidence<TimePoint> cachedFix;
cachedFix = getCachedFixes().floor(getDummyFixWithConfidence(timePoint));
WindWithConfidence<TimePoint> cachedFix = null;
WindWithConfidence<TimePoint> result = null;
final boolean nullResultCacheContains;
lockForRead();
try {
nullResultCacheContains = nullResultCacheContains(timePoint);
if (nullResultCacheContains) {
result = null;
} else {
synchronized (cache) {
cachedFix = cache.floor(getDummyFixWithConfidence(timePoint));
}
}
} finally {
unlockAfterRead();
}
if (!nullResultCacheContains) {
if (cachedFix == null || !cachedFix.getObject().getTimePoint().equals(timePoint)) {
result = getTrackedRace().getEstimatedWindDirectionWithConfidence(p, timePoint);
cache(timePoint, result);
@@ -403,7 +447,10 @@ public class TrackBasedEstimationWindTrackImpl extends VirtualWindTrackImpl impl
}
private boolean nullResultCacheContains(TimePoint timePoint) {
return timePointsWithCachedNullResultFastContains.contains(timePoint);
assertReadLock();
synchronized (timePointsWithCachedNullResultFastContains) {
return timePointsWithCachedNullResultFastContains.contains(timePoint);
}
}
@Override
@@ -503,29 +550,39 @@ public class TrackBasedEstimationWindTrackImpl extends VirtualWindTrackImpl impl
*/
@Override
protected WindWithConfidence<Pair<Position, TimePoint>> getAveragedWindUnsynchronized(Position p, TimePoint at) {
TimePoint floorTimePoint = virtualInternalRawFixes.floorToResolution(at);
TimePoint timePoint;
if (floorTimePoint.equals(at) ||
Math.abs(floorTimePoint.asMillis() - at.asMillis()) <
Math.abs(virtualInternalRawFixes.ceilingToResolution(at).asMillis() - at.asMillis())) {
timePoint = floorTimePoint;
} else {
timePoint = virtualInternalRawFixes.ceilingToResolution(at);
lockForRead();
try {
TimePoint floorTimePoint = virtualInternalRawFixes.floorToResolution(at);
TimePoint timePoint;
if (floorTimePoint.equals(at)
|| Math.abs(floorTimePoint.asMillis() - at.asMillis()) < Math.abs(virtualInternalRawFixes
.ceilingToResolution(at).asMillis() - at.asMillis())) {
timePoint = floorTimePoint;
} else {
timePoint = virtualInternalRawFixes.ceilingToResolution(at);
}
WindWithConfidence<TimePoint> preResult = virtualInternalRawFixes.getWindWithConfidence(p, timePoint);
// reduce confidence depending on how far *at* is away from the time point of the fix obtained
double confidenceMultiplier = weigher.getConfidence(timePoint, at);
WindWithConfidenceImpl<Pair<Position, TimePoint>> result = preResult == null ? null
: new WindWithConfidenceImpl<Pair<Position, TimePoint>>(preResult.getObject(), confidenceMultiplier
* preResult.getConfidence(),
/* relativeTo */new Pair<Position, TimePoint>(p, at), preResult.useSpeed());
return result;
} finally {
unlockAfterRead();
}
WindWithConfidence<TimePoint> preResult = virtualInternalRawFixes.getWindWithConfidence(p, timePoint);
// reduce confidence depending on how far *at* is away from the time point of the fix obtained
double confidenceMultiplier = weigher.getConfidence(timePoint, at);
WindWithConfidenceImpl<Pair<Position, TimePoint>> result = preResult == null ? null :
new WindWithConfidenceImpl<Pair<Position, TimePoint>>(
preResult.getObject(), confidenceMultiplier * preResult.getConfidence(),
/* relativeTo */ new Pair<Position, TimePoint>(p, at), preResult.useSpeed());
return result;
}
@Override
public String toString() {
return "This is the " + this.getClass().getName() + " object from " + virtualInternalRawFixes.getFrom()
+ " to " + virtualInternalRawFixes.getTo() + " for race " + getTrackedRace();
lockForRead();
try {
return "This is the " + this.getClass().getName() + " object from " + virtualInternalRawFixes.getFrom()
+ " to " + virtualInternalRawFixes.getTo() + " for race " + getTrackedRace();
} finally {
unlockAfterRead();
}
}
/**
@@ -5,6 +5,7 @@ import java.io.ObjectOutputStream;
import java.util.ConcurrentModificationException;
import java.util.Iterator;
import java.util.NavigableSet;
import java.util.concurrent.locks.ReentrantReadWriteLock;
import com.sap.sailing.domain.base.Timed;
import com.sap.sailing.domain.common.TimePoint;
@@ -18,7 +19,9 @@ public abstract class TrackImpl<FixType extends Timed> implements Track<FixType>
* The fixes, ordered by their time points
*/
private final NavigableSet<Timed> fixes;
private final ReentrantReadWriteLock readWriteLock;
protected static class DummyTimed implements Timed {
private static final long serialVersionUID = 6047311973718918856L;
private final TimePoint timePoint;
@@ -41,14 +44,38 @@ public abstract class TrackImpl<FixType extends Timed> implements Track<FixType>
}
protected TrackImpl(NavigableSet<Timed> fixes) {
this.readWriteLock = new ReentrantReadWriteLock();
this.fixes = fixes;
}
/**
* Synchronize the serialization such that no fixes are added while serializing
*/
private synchronized void writeObject(ObjectOutputStream s) throws IOException {
s.defaultWriteObject();
private void writeObject(ObjectOutputStream s) throws IOException {
lockForRead();
try {
s.defaultWriteObject();
} finally {
unlockAfterRead();
}
}
@Override
public void lockForRead() {
readWriteLock.readLock().lock();
}
@Override
public void unlockAfterRead() {
readWriteLock.readLock().unlock();
}
protected void lockForWrite() {
readWriteLock.writeLock().lock();
}
protected void unlockAfterWrite() {
readWriteLock.writeLock().unlock();
}
/**
@@ -60,7 +87,19 @@ public abstract class TrackImpl<FixType extends Timed> implements Track<FixType>
NavigableSet<FixType> result = (NavigableSet<FixType>) fixes;
return result;
}
protected void assertReadLock() {
if (readWriteLock.getReadHoldCount() < 1) {
throw new IllegalStateException("Caller must obtain read lock using lockForRead() before calling this method");
}
}
protected void assertWriteLock() {
if (readWriteLock.getWriteHoldCount() < 1) {
throw new IllegalStateException("Caller must obtain write lock using lockForWrite() before calling this method");
}
}
/**
* Callers that want to iterate over the collection returned need to synchronize on <code>this</code> object to
* avoid {@link ConcurrentModificationException}s.
@@ -80,6 +119,7 @@ public abstract class TrackImpl<FixType extends Timed> implements Track<FixType>
*/
@Override
public NavigableSet<FixType> getFixes() {
assertReadLock();
return new UnmodifiableNavigableSet<FixType>(getInternalFixes());
}
@@ -88,69 +128,121 @@ public abstract class TrackImpl<FixType extends Timed> implements Track<FixType>
*/
@Override
public NavigableSet<FixType> getRawFixes() {
assertReadLock();
return new UnmodifiableNavigableSet<FixType>(getInternalRawFixes());
}
@Override
public FixType getLastFixAtOrBefore(TimePoint timePoint) {
return (FixType) getInternalFixes().floor(getDummyFix(timePoint));
lockForRead();
try {
return (FixType) getInternalFixes().floor(getDummyFix(timePoint));
} finally {
unlockAfterRead();
}
}
@Override
public FixType getLastFixBefore(TimePoint timePoint) {
return (FixType) getInternalFixes().lower(getDummyFix(timePoint));
lockForRead();
try {
return (FixType) getInternalFixes().lower(getDummyFix(timePoint));
} finally {
unlockAfterRead();
}
}
@Override
public FixType getLastRawFixAtOrBefore(TimePoint timePoint) {
return (FixType) getInternalRawFixes().floor(getDummyFix(timePoint));
lockForRead();
try {
return (FixType) getInternalRawFixes().floor(getDummyFix(timePoint));
} finally {
unlockAfterRead();
}
}
@Override
public FixType getFirstRawFixAtOrAfter(TimePoint timePoint) {
return (FixType) getInternalRawFixes().ceiling(getDummyFix(timePoint));
lockForRead();
try {
return (FixType) getInternalRawFixes().ceiling(getDummyFix(timePoint));
} finally {
unlockAfterRead();
}
}
@Override
public FixType getFirstFixAtOrAfter(TimePoint timePoint) {
return (FixType) getInternalFixes().ceiling(getDummyFix(timePoint));
lockForRead();
try {
return (FixType) getInternalFixes().ceiling(getDummyFix(timePoint));
} finally {
unlockAfterRead();
}
}
@Override
public FixType getLastRawFixBefore(TimePoint timePoint) {
return (FixType) getInternalRawFixes().lower(getDummyFix(timePoint));
lockForRead();
try {
return (FixType) getInternalRawFixes().lower(getDummyFix(timePoint));
} finally {
unlockAfterRead();
}
}
@Override
public FixType getFirstFixAfter(TimePoint timePoint) {
return (FixType) getInternalFixes().higher(getDummyFix(timePoint));
lockForRead();
try {
return (FixType) getInternalFixes().higher(getDummyFix(timePoint));
} finally {
unlockAfterRead();
}
}
@Override
public FixType getFirstRawFixAfter(TimePoint timePoint) {
return (FixType) getInternalRawFixes().higher(getDummyFix(timePoint));
lockForRead();
try {
return (FixType) getInternalRawFixes().higher(getDummyFix(timePoint));
} finally {
unlockAfterRead();
}
}
@Override
public FixType getFirstRawFix() {
if (getInternalFixes().isEmpty()) {
return null;
} else {
return (FixType) getInternalFixes().first();
lockForRead();
try {
if (getInternalFixes().isEmpty()) {
return null;
} else {
return (FixType) getInternalFixes().first();
}
} finally {
unlockAfterRead();
}
}
@Override
public FixType getLastRawFix() {
if (getInternalRawFixes().isEmpty()) {
return null;
} else {
return (FixType) getInternalRawFixes().last();
lockForRead();
try {
if (getInternalRawFixes().isEmpty()) {
return null;
} else {
return (FixType) getInternalRawFixes().last();
}
} finally {
unlockAfterRead();
}
}
@Override
public Iterator<FixType> getFixesIterator(TimePoint startingAt, boolean inclusive) {
assertReadLock();
Iterator<FixType> result = (Iterator<FixType>) getInternalFixes().tailSet(
getDummyFix(startingAt), inclusive).iterator();
return result;
@@ -171,6 +263,7 @@ public abstract class TrackImpl<FixType extends Timed> implements Track<FixType>
@Override
public Iterator<FixType> getRawFixesIterator(TimePoint startingAt, boolean inclusive) {
assertReadLock();
Iterator<FixType> result = (Iterator<FixType>) getInternalRawFixes().tailSet(
getDummyFix(startingAt), inclusive).iterator();
return result;
@@ -18,6 +18,9 @@ import java.util.Set;
import java.util.Timer;
import java.util.TimerTask;
import java.util.concurrent.ConcurrentSkipListSet;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReadWriteLock;
import java.util.concurrent.locks.ReentrantReadWriteLock;
import java.util.logging.Level;
import java.util.logging.Logger;
@@ -144,7 +147,15 @@ public abstract class TrackedRaceImpl implements TrackedRace, CourseListener {
private long updateCount;
private final Map<TimePoint, List<Competitor>> competitorRankings;
private transient Map<TimePoint, List<Competitor>> competitorRankings;
/**
* The locks managed here correspond with the {@link #competitorRankings} structure. When
* {@link #getCompetitorsFromBestToWorst(TimePoint)} starts to compute rankings, it locks the write lock for the
* time point. Readers use the read lock. Checking / entering a lock into this map uses <code>synchronized</code> on
* the map itself.
*/
private transient Map<TimePoint, ReadWriteLock> competitorRankingsLocks;
/**
* legs appear in the order in which they appear in the race's course
@@ -256,6 +267,7 @@ public abstract class TrackedRaceImpl implements TrackedRace, CourseListener {
windTracks.put(trackBasedWindSource, getOrCreateWindTrack(trackBasedWindSource, delayForWindEstimationCacheInvalidation));
this.trackedRegatta = trackedRegatta;
competitorRankings = new HashMap<TimePoint, List<Competitor>>();
competitorRankingsLocks = new HashMap<TimePoint, ReadWriteLock>();
}
/**
@@ -265,6 +277,8 @@ public abstract class TrackedRaceImpl implements TrackedRace, CourseListener {
private void readObject(ObjectInputStream ois) throws ClassNotFoundException, IOException {
ois.defaultReadObject();
windStore = EmptyWindStore.INSTANCE;
competitorRankings = new HashMap<TimePoint, List<Competitor>>();
competitorRankingsLocks = new HashMap<TimePoint, ReadWriteLock>();
}
/**
@@ -671,8 +685,26 @@ public abstract class TrackedRaceImpl implements TrackedRace, CourseListener {
@Override
public List<Competitor> getCompetitorsFromBestToWorst(TimePoint timePoint) {
ReadWriteLock readWriteLock;
synchronized (competitorRankingsLocks) {
readWriteLock = competitorRankingsLocks.get(timePoint);
if (readWriteLock == null) {
readWriteLock = new ReentrantReadWriteLock();
competitorRankingsLocks.put(timePoint, readWriteLock);
}
}
List<Competitor> rankedCompetitors;
final Lock lock;
synchronized (competitorRankings) {
List<Competitor> rankedCompetitors = competitorRankings.get(timePoint);
rankedCompetitors = competitorRankings.get(timePoint);
if (rankedCompetitors == null) {
lock = readWriteLock.writeLock();
} else {
lock = readWriteLock.readLock();
}
}
lock.lock();
try {
if (rankedCompetitors == null) {
RaceRankComparator comparator = new RaceRankComparator(this, timePoint);
rankedCompetitors = new ArrayList<Competitor>();
@@ -680,11 +712,14 @@ public abstract class TrackedRaceImpl implements TrackedRace, CourseListener {
rankedCompetitors.add(c);
}
Collections.sort(rankedCompetitors, comparator);
competitorRankings.put(timePoint, rankedCompetitors);
synchronized (competitorRankings) {
competitorRankings.put(timePoint, rankedCompetitors);
}
}
return rankedCompetitors;
} finally {
lock.unlock();
}
}
@Override
@@ -1013,6 +1048,9 @@ public abstract class TrackedRaceImpl implements TrackedRace, CourseListener {
synchronized (competitorRankings) {
competitorRankings.clear();
}
synchronized (competitorRankingsLocks) {
competitorRankingsLocks.clear();
}
synchronized (maneuverCache) {
maneuverCache.clear();
}
@@ -99,8 +99,11 @@ public class WindTrackImpl extends TrackImpl<Wind> implements WindTrack {
@Override
public void add(Wind wind) {
CompactWindImpl compactWind = new CompactWindImpl(wind);
synchronized (this) {
lockForWrite();
try {
getInternalRawFixes().add(compactWind);
} finally {
unlockAfterWrite();
}
notifyListenersAboutReceive(compactWind);
}
@@ -175,14 +178,15 @@ public class WindTrackImpl extends TrackImpl<Wind> implements WindTrack {
* <code>p</code> is used as the result's position and may be used for confidence determination.
*/
protected WindWithConfidence<Pair<Position, TimePoint>> getAveragedWindUnsynchronized(Position p, TimePoint at) {
List<WindWithConfidence<Pair<Position, TimePoint>>> windFixesToAverage = new ArrayList<WindWithConfidence<Pair<Position, TimePoint>>>();
// don't measure speed with separate confidence; return confidence obtained from averaging bearings
ConfidenceBasedWindAverager<Pair<Position, TimePoint>> windAverager = ConfidenceFactory.INSTANCE
.createWindAverager(new PositionAndTimePointWeigher(
/* halfConfidenceAfterMilliseconds */getMillisecondsOverWhichToAverageWind() / 10));
DummyWind atTimed = new DummyWind(at);
Pair<Position, TimePoint> relativeTo = new Pair<Position, TimePoint>(p, at);
synchronized (this) {
lockForRead();
try {
List<WindWithConfidence<Pair<Position, TimePoint>>> windFixesToAverage = new ArrayList<WindWithConfidence<Pair<Position, TimePoint>>>();
// don't measure speed with separate confidence; return confidence obtained from averaging bearings
ConfidenceBasedWindAverager<Pair<Position, TimePoint>> windAverager = ConfidenceFactory.INSTANCE
.createWindAverager(new PositionAndTimePointWeigher(
/* halfConfidenceAfterMilliseconds */getMillisecondsOverWhichToAverageWind() / 10));
DummyWind atTimed = new DummyWind(at);
Pair<Position, TimePoint> relativeTo = new Pair<Position, TimePoint>(p, at);
NavigableSet<Wind> beforeSet = getInternalFixes().headSet(atTimed, /* inclusive */false);
NavigableSet<Wind> afterSet = getInternalFixes().tailSet(atTimed, /* inclusive */true);
Iterator<Wind> beforeIter = beforeSet.descendingIterator();
@@ -235,12 +239,15 @@ public class WindTrackImpl extends TrackImpl<Wind> implements WindTrack {
}
} while (beforeIntervalLength + afterIntervalLength < getMillisecondsOverWhichToAverageWind()
&& (beforeWind != null || afterWind != null));
}
if (windFixesToAverage.isEmpty()) {
return null;
} else {
WindWithConfidence<Pair<Position, TimePoint>> average = windAverager.getAverage(windFixesToAverage, relativeTo);
return average;
if (windFixesToAverage.isEmpty()) {
return null;
} else {
WindWithConfidence<Pair<Position, TimePoint>> average = windAverager.getAverage(windFixesToAverage,
relativeTo);
return average;
}
} finally {
unlockAfterRead();
}
}
@@ -254,35 +261,45 @@ public class WindTrackImpl extends TrackImpl<Wind> implements WindTrack {
@Override
public String toString() {
StringBuilder result = new StringBuilder();
synchronized (this) {
for (Wind wind : getRawFixes()) {
result.append(wind);
result.append(" avg(");
result.append(getMillisecondsOverWhichToAverageWind());
if (wind == null) {
result.append("ms)");
} else {
result.append("ms): ");
result.append(getAveragedWind(wind.getPosition(), wind.getTimePoint()));
lockForRead();
try {
StringBuilder result = new StringBuilder();
synchronized (this) {
for (Wind wind : getRawFixes()) {
result.append(wind);
result.append(" avg(");
result.append(getMillisecondsOverWhichToAverageWind());
if (wind == null) {
result.append("ms)");
} else {
result.append("ms): ");
result.append(getAveragedWind(wind.getPosition(), wind.getTimePoint()));
}
result.append("\n");
}
result.append("\n");
}
return result.toString();
} finally {
unlockAfterRead();
}
return result.toString();
}
public String toCSV() {
StringBuilder result = new StringBuilder();
synchronized (this) {
for (Wind wind : getRawFixes()) {
append(result, wind);
Wind estimate = getAveragedWind(wind.getPosition(), wind.getTimePoint());
append(result, estimate);
result.append("\n");
lockForRead();
try {
StringBuilder result = new StringBuilder();
synchronized (this) {
for (Wind wind : getRawFixes()) {
append(result, wind);
Wind estimate = getAveragedWind(wind.getPosition(), wind.getTimePoint());
append(result, estimate);
result.append("\n");
}
}
return result.toString();
} finally {
unlockAfterRead();
}
return result.toString();
}
private void append(StringBuilder result, Wind wind) {
@@ -358,7 +375,12 @@ public class WindTrackImpl extends TrackImpl<Wind> implements WindTrack {
@Override
public void remove(Wind wind) {
getInternalRawFixes().remove(wind);
lockForWrite();
try {
getInternalRawFixes().remove(wind);
} finally {
unlockAfterWrite();
}
notifyListenersAboutRemoval(wind);
}
@@ -1390,7 +1390,9 @@ public class SailingServiceImpl extends RemoteServiceServlet implements SailingS
for (Buoy startBuoy : waypoint.getBuoys()) {
final Position estimatedBuoyPosition = trackedRace.getOrCreateTrack(startBuoy)
.getEstimatedPosition(timePoint, /* extrapolate */false);
startBuoyPositions.add(new PositionDTO(estimatedBuoyPosition.getLatDeg(), estimatedBuoyPosition.getLngDeg()));
if (estimatedBuoyPosition != null) {
startBuoyPositions.add(new PositionDTO(estimatedBuoyPosition.getLatDeg(), estimatedBuoyPosition.getLngDeg()));
}
}
return startBuoyPositions;
}