mirror of
https://github.com/eclipse-sailing-analytics/sailing-analytics.git
synced 2026-09-25 06:58:39 +00:00
Bug 4489: Initial implementation to ensure replication state for out of
thread event processing in TrackedRegattaImpl
This commit is contained in:
+2
-1
@@ -11,6 +11,7 @@ import java.io.IOException;
|
||||
import java.net.MalformedURLException;
|
||||
import java.net.URISyntaxException;
|
||||
import java.util.Collections;
|
||||
import java.util.Optional;
|
||||
import java.util.UUID;
|
||||
|
||||
import org.junit.After;
|
||||
@@ -228,7 +229,7 @@ public class CreateAndTrackWithRaceLogTest extends RaceLogTrackingTestHelper {
|
||||
public void raceAdded(TrackedRace trackedRace) {
|
||||
}
|
||||
};
|
||||
raceHandle.getTrackedRegatta().addRaceListener(raceListener);
|
||||
raceHandle.getTrackedRegatta().addRaceListener(raceListener, Optional.empty());
|
||||
raceHandle.getTrackedRegatta().removeRaceListener(raceListener).get();
|
||||
}
|
||||
|
||||
|
||||
+3
-1
@@ -7,6 +7,7 @@ import java.util.Collections;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.Map.Entry;
|
||||
import java.util.logging.Level;
|
||||
import java.util.logging.Logger;
|
||||
@@ -304,7 +305,8 @@ public class RaceLogRaceTracker extends AbstractRaceTrackerBaseImpl {
|
||||
raceColumn.setRaceIdentifier(fleet, trackedRegatta.getRegatta().getRaceIdentifier(raceDef));
|
||||
trackedRace = trackedRegatta.createTrackedRace(raceDef, sidelines, windStore,
|
||||
params.getDelayToLiveInMillis(), WindTrack.DEFAULT_MILLISECONDS_OVER_WHICH_TO_AVERAGE_WIND,
|
||||
boatClass.getApproximateManeuverDurationInMilliseconds(), null, /*useMarkPassingCalculator*/ true, raceLogResolver);
|
||||
boatClass.getApproximateManeuverDurationInMilliseconds(), null, /*useMarkPassingCalculator*/ true, raceLogResolver,
|
||||
/* Not needed because the RaceTracker is not active on a replica */ Optional.empty());
|
||||
notifyRaceCreationListeners();
|
||||
logger.info(String.format("Started tracking race-log race (%s)", raceLog));
|
||||
// this wakes up all waiting race handles
|
||||
|
||||
+3
-1
@@ -10,6 +10,7 @@ import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.NavigableSet;
|
||||
import java.util.Optional;
|
||||
import java.util.TreeMap;
|
||||
import java.util.logging.Level;
|
||||
import java.util.logging.Logger;
|
||||
@@ -487,7 +488,8 @@ public class SwissTimingRaceTrackerImpl extends AbstractRaceTrackerImpl
|
||||
// we already know our single RaceDefinition
|
||||
assert SwissTimingRaceTrackerImpl.this.race == race;
|
||||
}
|
||||
}, useInternalMarkPassingAlgorithm, raceLogResolver);
|
||||
}, useInternalMarkPassingAlgorithm, raceLogResolver,
|
||||
/* Not needed because the RaceTracker is not active on a replica */ Optional.empty());
|
||||
notifyRaceCreationListeners();
|
||||
logger.info("Created SwissTiming RaceDefinition and TrackedRace for "+race.getName());
|
||||
}
|
||||
|
||||
+3
-1
@@ -11,6 +11,7 @@ import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Map.Entry;
|
||||
import java.util.NavigableSet;
|
||||
import java.util.Optional;
|
||||
import java.util.Set;
|
||||
import java.util.TimeZone;
|
||||
import java.util.logging.Logger;
|
||||
@@ -400,7 +401,8 @@ public class SwissTimingReplayToDomainAdapter extends SwissTimingReplayAdapter i
|
||||
TrackedRace.DEFAULT_LIVE_DELAY_IN_MILLISECONDS,
|
||||
WindTrack.DEFAULT_MILLISECONDS_OVER_WHICH_TO_AVERAGE_WIND,
|
||||
/* time over which to average speed: */ race.getBoatClass().getApproximateManeuverDurationInMilliseconds(),
|
||||
/* raceDefinitionSetToUpdate */ null, useInternalMarkPassingAlgorithm, raceLogResolver);
|
||||
/* raceDefinitionSetToUpdate */ null, useInternalMarkPassingAlgorithm, raceLogResolver,
|
||||
/* Not needed because the RaceTracker is not active on a replica */ Optional.empty());
|
||||
trackedRace.onStatusChanged(this, new TrackedRaceStatusImpl(TrackedRaceStatusEnum.LOADING, 0));
|
||||
TimePoint bestStartTimeKnownSoFar = bestStartTimePerRaceID.get(currentRaceID);
|
||||
if (bestStartTimeKnownSoFar != null) {
|
||||
|
||||
+2
-1
@@ -6,6 +6,7 @@ import java.net.MalformedURLException;
|
||||
import java.net.URISyntaxException;
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
|
||||
import org.junit.Ignore;
|
||||
import org.junit.Test;
|
||||
@@ -74,7 +75,7 @@ public class FetchTracksAndStoreLocallyTest extends OnlineTracTracBasedTest {
|
||||
@Override
|
||||
public void raceRemoved(TrackedRace trackedRace) {
|
||||
}
|
||||
});
|
||||
}, Optional.empty());
|
||||
super.completeSetupLaunchingControllerAndWaitForRaceDefinition(ReceiverType.RACECOURSE,
|
||||
ReceiverType.RACESTARTFINISH, ReceiverType.RAWPOSITIONS);
|
||||
}
|
||||
|
||||
+2
-1
@@ -5,6 +5,7 @@ import static org.mockito.Mockito.mock;
|
||||
|
||||
import java.net.MalformedURLException;
|
||||
import java.net.URISyntaxException;
|
||||
import java.util.Optional;
|
||||
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
@@ -77,7 +78,7 @@ public class ReceiveTrackingDataTest extends AbstractTracTracLiveTest {
|
||||
@Override
|
||||
public void raceRemoved(TrackedRace trackedRace) {
|
||||
}
|
||||
});
|
||||
}, Optional.empty());
|
||||
for (Receiver receiver : domainFactory
|
||||
.getUpdateReceivers(trackedRegatta, /* delayToLiveInMillis */0l,
|
||||
/* simulator */null, EmptyWindStore.INSTANCE, new DynamicRaceDefinitionSet() {
|
||||
|
||||
+7
-4
@@ -5,6 +5,7 @@ import static org.mockito.Mockito.when;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.net.MalformedURLException;
|
||||
import java.util.Optional;
|
||||
import java.util.UUID;
|
||||
import java.util.concurrent.BrokenBarrierException;
|
||||
import java.util.concurrent.CyclicBarrier;
|
||||
@@ -33,6 +34,7 @@ import com.sap.sailing.domain.tracking.impl.DynamicTrackedRaceImpl;
|
||||
import com.sap.sailing.domain.tracking.impl.DynamicTrackedRegattaImpl;
|
||||
import com.sap.sailing.server.RacingEventService;
|
||||
import com.sap.sailing.server.impl.RacingEventServiceImpl;
|
||||
import com.sap.sse.util.ThreadLocalTransporter;
|
||||
|
||||
public class TestDeadlockInRegattaListener {
|
||||
@Rule
|
||||
@@ -56,13 +58,14 @@ public class TestDeadlockInRegattaListener {
|
||||
private static final long serialVersionUID = -3599667964201700780L;
|
||||
|
||||
@Override
|
||||
protected void notifyListenersAboutTrackedRaceRemoved(TrackedRace trackedRace) {
|
||||
protected void notifyListenersAboutTrackedRaceRemoved(TrackedRace trackedRace,
|
||||
Optional<ThreadLocalTransporter> threadLocalTransporter) {
|
||||
try {
|
||||
latch.await();
|
||||
} catch (InterruptedException | BrokenBarrierException e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
super.notifyListenersAboutTrackedRaceRemoved(trackedRace);
|
||||
super.notifyListenersAboutTrackedRaceRemoved(trackedRace, Optional.empty());
|
||||
}
|
||||
};
|
||||
RacingEventServiceImpl racingEventService = new RacingEventServiceImpl() {
|
||||
@@ -125,10 +128,10 @@ public class TestDeadlockInRegattaListener {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
}).start();
|
||||
trackedRegatta.addTrackedRace(trackedRace1);
|
||||
trackedRegatta.addTrackedRace(trackedRace1, Optional.empty());
|
||||
// the following runs into RacingEventService.getRaceTrackerByRegattaAndRaceIdentifier
|
||||
// which waits for the latch based on the override above while in synchronized RegattaListener.raceAdded
|
||||
new Thread(()->trackedRegatta.addTrackedRace(trackedRace2)).start();
|
||||
new Thread(()->trackedRegatta.addTrackedRace(trackedRace2, Optional.empty())).start();
|
||||
monitorOnRegattaListenerLatch.await();
|
||||
// the following awaits the latch in TrackedRegattaImpl.notifyListenersAboutTrackedRaceRemoved
|
||||
// after the write lock has been obtained but before the synchronized RegattaListener.raceRemoved method
|
||||
|
||||
+6
-4
@@ -4,6 +4,7 @@ import java.io.Serializable;
|
||||
import java.util.Collections;
|
||||
import java.util.List;
|
||||
import java.util.NavigableSet;
|
||||
import java.util.Optional;
|
||||
import java.util.Set;
|
||||
import java.util.TreeSet;
|
||||
import java.util.concurrent.Future;
|
||||
@@ -89,6 +90,7 @@ import com.sap.sse.common.IsManagedByCache;
|
||||
import com.sap.sse.common.TimePoint;
|
||||
import com.sap.sse.common.Util;
|
||||
import com.sap.sse.common.Util.Pair;
|
||||
import com.sap.sse.util.ThreadLocalTransporter;
|
||||
|
||||
public class MockedTrackedRace implements DynamicTrackedRace {
|
||||
private static final long serialVersionUID = 5827912985564121181L;
|
||||
@@ -565,15 +567,15 @@ public class MockedTrackedRace implements DynamicTrackedRace {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addTrackedRace(TrackedRace trackedRace) {
|
||||
public void addTrackedRace(TrackedRace trackedRace, Optional<ThreadLocalTransporter> threadLocalTransporter) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void removeTrackedRace(TrackedRace trackedRace) {
|
||||
public void removeTrackedRace(TrackedRace trackedRace, Optional<ThreadLocalTransporter> threadLocalTransporter) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addRaceListener(RaceListener listener) {
|
||||
public void addRaceListener(RaceListener listener, Optional<ThreadLocalTransporter> threadLocalTransporter) {
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -596,7 +598,7 @@ public class MockedTrackedRace implements DynamicTrackedRace {
|
||||
WindStore windStore, long delayToLiveInMillis,
|
||||
long millisecondsOverWhichToAverageWind, long millisecondsOverWhichToAverageSpeed,
|
||||
DynamicRaceDefinitionSet raceDefinitionSetToUpdate, boolean useMarkPassingcalculator,
|
||||
RaceLogResolver raceLogResolver) {
|
||||
RaceLogResolver raceLogResolver, Optional<ThreadLocalTransporter> threadLocalTransporter) {
|
||||
return null;
|
||||
}
|
||||
|
||||
|
||||
+4
-3
@@ -4,6 +4,7 @@ import static org.mockito.Mockito.mock;
|
||||
|
||||
import java.util.Arrays;
|
||||
import java.util.Collections;
|
||||
import java.util.Optional;
|
||||
import java.util.concurrent.CyclicBarrier;
|
||||
import java.util.concurrent.Phaser;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
@@ -74,11 +75,11 @@ public class TrackedRegattaTest {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
}
|
||||
});
|
||||
}, Optional.empty());
|
||||
|
||||
DynamicTrackedRace race1 = createRace("R1");
|
||||
Thread thread1 = new Thread(() -> {
|
||||
regatta.addTrackedRace(race1);
|
||||
regatta.addTrackedRace(race1, Optional.empty());
|
||||
});
|
||||
thread1.start();
|
||||
// This ensures, that the add event is being processed but is not finished because
|
||||
@@ -88,7 +89,7 @@ public class TrackedRegattaTest {
|
||||
addPhaser.arriveAndAwaitAdvance();
|
||||
|
||||
Thread thread2 = new Thread(() -> {
|
||||
regatta.removeTrackedRace(race1);
|
||||
regatta.removeTrackedRace(race1, Optional.empty());
|
||||
});
|
||||
thread2.start();
|
||||
// If the implementation ensures that the events are fired in order,
|
||||
|
||||
+2
-1
@@ -13,6 +13,7 @@ import java.util.HashSet;
|
||||
import java.util.Iterator;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.Map.Entry;
|
||||
import java.util.Set;
|
||||
import java.util.UUID;
|
||||
@@ -616,7 +617,7 @@ public class DomainFactoryImpl implements DomainFactory {
|
||||
return trackedRegatta.createTrackedRace(race, sidelines,
|
||||
windStore, delayToLiveInMillis, millisecondsOverWhichToAverageWind,
|
||||
/* time over which to average speed: */ race.getBoatClass().getApproximateManeuverDurationInMilliseconds(),
|
||||
raceDefinitionSetToUpdate, useMarkPassingCalculator, raceLogResolver);
|
||||
raceDefinitionSetToUpdate, useMarkPassingCalculator, raceLogResolver, Optional.empty());
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
+3
-1
@@ -5,6 +5,7 @@ import java.util.ArrayList;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.logging.Level;
|
||||
import java.util.logging.Logger;
|
||||
@@ -241,7 +242,8 @@ public class RaceCourseReceiver extends AbstractReceiverWithQueue<IControlRoute,
|
||||
DynamicTrackedRace trackedRace = getTrackedRegatta().createTrackedRace(race, sidelines,
|
||||
windStore, delayToLiveInMillis, millisecondsOverWhichToAverageWind,
|
||||
/* time over which to average speed: */ race.getBoatClass().getApproximateManeuverDurationInMilliseconds(),
|
||||
raceDefinitionSetToUpdate, useInternalMarkPassingAlgorithm, raceLogResolver);
|
||||
raceDefinitionSetToUpdate, useInternalMarkPassingAlgorithm, raceLogResolver,
|
||||
/* Not needed because the RaceTracker is not active on a replica */ Optional.empty());
|
||||
if (runAfterCreatingTrackedRace != null) {
|
||||
runAfterCreatingTrackedRace.accept(trackedRace);
|
||||
}
|
||||
|
||||
+2
-1
@@ -4,6 +4,7 @@ import java.util.ArrayList;
|
||||
import java.util.Iterator;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.Timer;
|
||||
import java.util.TimerTask;
|
||||
import java.util.logging.Logger;
|
||||
@@ -70,7 +71,7 @@ public class Simulator {
|
||||
stop(); // stop simulator when tracked race is removed from its regatta
|
||||
}
|
||||
}
|
||||
});
|
||||
}, /* No replication handling necessary */ Optional.empty());
|
||||
startWindPlayer();
|
||||
}
|
||||
|
||||
|
||||
+2
-1
@@ -2,6 +2,7 @@ package com.sap.sailing.domain.tracking;
|
||||
|
||||
import java.io.Serializable;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import java.util.logging.Level;
|
||||
@@ -107,7 +108,7 @@ public abstract class AbstractTrackedRegattaAndRaceObserver implements TrackedRe
|
||||
RegattaListener.this.raceAdded(trackedRace);
|
||||
}
|
||||
};
|
||||
trackedRegatta.addRaceListener(raceListener);
|
||||
trackedRegatta.addRaceListener(raceListener, /* Not replicated */ Optional.empty());
|
||||
}
|
||||
|
||||
public synchronized void raceRemoved(TrackedRace trackedRace) {
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package com.sap.sailing.domain.tracking;
|
||||
|
||||
import java.io.Serializable;
|
||||
import java.util.Optional;
|
||||
import java.util.concurrent.Future;
|
||||
|
||||
import com.sap.sailing.domain.abstractlog.race.analyzing.impl.RaceLogResolver;
|
||||
@@ -11,6 +12,7 @@ import com.sap.sailing.domain.base.Sideline;
|
||||
import com.sap.sailing.domain.base.impl.TrackedRaces;
|
||||
import com.sap.sailing.domain.common.NoWindException;
|
||||
import com.sap.sse.common.TimePoint;
|
||||
import com.sap.sse.util.ThreadLocalTransporter;
|
||||
|
||||
/**
|
||||
* Manages a set of {@link TrackedRace} objects that belong to the same {@link Regatta} (regatta, sailing regatta for a
|
||||
@@ -65,7 +67,8 @@ public interface TrackedRegatta extends Serializable {
|
||||
*/
|
||||
DynamicTrackedRace createTrackedRace(RaceDefinition raceDefinition, Iterable<Sideline> sidelines, WindStore windStore,
|
||||
long delayToLiveInMillis, long millisecondsOverWhichToAverageWind, long millisecondsOverWhichToAverageSpeed,
|
||||
DynamicRaceDefinitionSet raceDefinitionSetToUpdate, boolean useInternalMarkPassingAlgorithm, RaceLogResolver raceLogResolver);
|
||||
DynamicRaceDefinitionSet raceDefinitionSetToUpdate, boolean useInternalMarkPassingAlgorithm, RaceLogResolver raceLogResolver,
|
||||
Optional<ThreadLocalTransporter> beforeAndAfterNotificationHandler);
|
||||
|
||||
/**
|
||||
* Obtains the tracked race for <code>race</code>. Blocks until the tracked race has been created
|
||||
@@ -79,16 +82,16 @@ public interface TrackedRegatta extends Serializable {
|
||||
*/
|
||||
TrackedRace getExistingTrackedRace(RaceDefinition race);
|
||||
|
||||
void addTrackedRace(TrackedRace trackedRace);
|
||||
void addTrackedRace(TrackedRace trackedRace, Optional<ThreadLocalTransporter> beforeAndAfterNotificationHandler);
|
||||
|
||||
void removeTrackedRace(TrackedRace trackedRace);
|
||||
void removeTrackedRace(TrackedRace trackedRace, Optional<ThreadLocalTransporter> beforeAndAfterNotificationHandler);
|
||||
|
||||
/**
|
||||
* Listener will be notified when {@link #addTrackedRace(TrackedRace)} is called and
|
||||
* upon registration for each tracked race already known. Therefore, the listener
|
||||
* won't miss any tracked race.
|
||||
*/
|
||||
void addRaceListener(RaceListener listener);
|
||||
void addRaceListener(RaceListener listener, Optional<ThreadLocalTransporter> beforeAndAfterNotificationHandler);
|
||||
|
||||
/**
|
||||
* Removes the given listener and returns a {@link Future} that will be completed
|
||||
|
||||
+6
-2
@@ -1,5 +1,7 @@
|
||||
package com.sap.sailing.domain.tracking.impl;
|
||||
|
||||
import java.util.Optional;
|
||||
|
||||
import com.sap.sailing.domain.abstractlog.race.analyzing.impl.RaceLogResolver;
|
||||
import com.sap.sailing.domain.base.RaceDefinition;
|
||||
import com.sap.sailing.domain.base.Regatta;
|
||||
@@ -8,6 +10,7 @@ import com.sap.sailing.domain.tracking.DynamicRaceDefinitionSet;
|
||||
import com.sap.sailing.domain.tracking.DynamicTrackedRace;
|
||||
import com.sap.sailing.domain.tracking.DynamicTrackedRegatta;
|
||||
import com.sap.sailing.domain.tracking.WindStore;
|
||||
import com.sap.sse.util.ThreadLocalTransporter;
|
||||
|
||||
public class DynamicTrackedRegattaImpl extends TrackedRegattaImpl implements DynamicTrackedRegatta {
|
||||
private static final long serialVersionUID = -90155868534737120L;
|
||||
@@ -35,9 +38,10 @@ public class DynamicTrackedRegattaImpl extends TrackedRegattaImpl implements Dyn
|
||||
@Override
|
||||
public DynamicTrackedRace createTrackedRace(RaceDefinition raceDefinition, Iterable<Sideline> sidelines, WindStore windStore,
|
||||
long delayToLiveInMillis, long millisecondsOverWhichToAverageWind, long millisecondsOverWhichToAverageSpeed,
|
||||
DynamicRaceDefinitionSet raceDefinitionSetToUpdate, boolean useMarkPassingCalculator, RaceLogResolver raceLogResolver) {
|
||||
DynamicRaceDefinitionSet raceDefinitionSetToUpdate, boolean useMarkPassingCalculator, RaceLogResolver raceLogResolver,
|
||||
Optional<ThreadLocalTransporter> threadLocalTransporter) {
|
||||
return (DynamicTrackedRace) super.createTrackedRace(raceDefinition, sidelines, windStore, delayToLiveInMillis,
|
||||
millisecondsOverWhichToAverageWind,
|
||||
millisecondsOverWhichToAverageSpeed, raceDefinitionSetToUpdate, useMarkPassingCalculator, raceLogResolver);
|
||||
millisecondsOverWhichToAverageSpeed, raceDefinitionSetToUpdate, useMarkPassingCalculator, raceLogResolver, threadLocalTransporter);
|
||||
}
|
||||
}
|
||||
|
||||
+1
-1
@@ -548,7 +548,7 @@ public abstract class TrackedRaceImpl extends TrackedRaceWithWindEssentials impl
|
||||
markPassingCalculator.stop();
|
||||
}
|
||||
}
|
||||
});
|
||||
}, /* Not relevant For replication */ Optional.empty());
|
||||
} else {
|
||||
markPassingCalculator = null;
|
||||
}
|
||||
|
||||
+37
-19
@@ -9,6 +9,7 @@ import java.util.HashMap;
|
||||
import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
@@ -34,6 +35,7 @@ import com.sap.sse.common.TimePoint;
|
||||
import com.sap.sse.common.Util;
|
||||
import com.sap.sse.concurrent.LockUtil;
|
||||
import com.sap.sse.concurrent.NamedReentrantReadWriteLock;
|
||||
import com.sap.sse.util.ThreadLocalTransporter;
|
||||
|
||||
public class TrackedRegattaImpl implements TrackedRegatta {
|
||||
private static final long serialVersionUID = 6480508193567014285L;
|
||||
@@ -128,7 +130,7 @@ public class TrackedRegattaImpl implements TrackedRegatta {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addTrackedRace(TrackedRace trackedRace) {
|
||||
public void addTrackedRace(TrackedRace trackedRace, Optional<ThreadLocalTransporter> threadLocalTransporter) {
|
||||
final TrackedRace oldTrackedRace;
|
||||
lockTrackedRacesForWrite();
|
||||
try {
|
||||
@@ -136,39 +138,51 @@ public class TrackedRegattaImpl implements TrackedRegatta {
|
||||
" with regatta hash code "+getRegatta().hashCode());
|
||||
oldTrackedRace = trackedRaces.put(trackedRace.getRace(), trackedRace);
|
||||
if (oldTrackedRace != trackedRace) {
|
||||
notifyListenersAboutTrackedRaceAdded(trackedRace);
|
||||
notifyListenersAboutTrackedRaceAdded(trackedRace, threadLocalTransporter);
|
||||
}
|
||||
} finally {
|
||||
unlockTrackedRacesAfterWrite();
|
||||
}
|
||||
}
|
||||
|
||||
protected void notifyListenersAboutTrackedRaceAdded(TrackedRace trackedRace) {
|
||||
enqueEvent(listener -> listener.raceAdded(trackedRace));
|
||||
protected void notifyListenersAboutTrackedRaceAdded(TrackedRace trackedRace, Optional<ThreadLocalTransporter> threadLocalTransporter) {
|
||||
enqueEvent(listener -> listener.raceAdded(trackedRace), threadLocalTransporter);
|
||||
}
|
||||
|
||||
protected void enqueEvent(Consumer<RaceListener> fireEventCallback) {
|
||||
protected void enqueEvent(Consumer<RaceListener> fireEventCallback, Optional<ThreadLocalTransporter> threadLocalTransporter) {
|
||||
final Set<RaceListener> listenersToInform = new HashSet<>(raceListeners.keySet());
|
||||
threadLocalTransporter.ifPresent(ThreadLocalTransporter::rememberThreadLocalStates);
|
||||
eventQueue.addWork(() -> {
|
||||
for (RaceListener listener : listenersToInform) {
|
||||
fireEventCallback.accept(listener);
|
||||
}
|
||||
withBeforeAndAfterHandling(threadLocalTransporter, () -> {
|
||||
for (RaceListener listener : listenersToInform) {
|
||||
fireEventCallback.accept(listener);
|
||||
}
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
private void withBeforeAndAfterHandling(Optional<ThreadLocalTransporter> threadLocalTransporter, Runnable action) {
|
||||
threadLocalTransporter.ifPresent(ThreadLocalTransporter::pushThreadLocalStates);
|
||||
try {
|
||||
action.run();
|
||||
} finally {
|
||||
threadLocalTransporter.ifPresent(ThreadLocalTransporter::popThreadLocalStates);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void removeTrackedRace(TrackedRace trackedRace) {
|
||||
public void removeTrackedRace(TrackedRace trackedRace, Optional<ThreadLocalTransporter> threadLocalTransporter) {
|
||||
lockTrackedRacesForWrite();
|
||||
try {
|
||||
trackedRaces.remove(trackedRace.getRace());
|
||||
notifyListenersAboutTrackedRaceRemoved(trackedRace);
|
||||
notifyListenersAboutTrackedRaceRemoved(trackedRace, threadLocalTransporter);
|
||||
} finally {
|
||||
unlockTrackedRacesAfterWrite();
|
||||
}
|
||||
}
|
||||
|
||||
protected void notifyListenersAboutTrackedRaceRemoved(TrackedRace trackedRace) {
|
||||
enqueEvent(listener -> listener.raceRemoved(trackedRace));
|
||||
protected void notifyListenersAboutTrackedRaceRemoved(TrackedRace trackedRace, Optional<ThreadLocalTransporter> threadLocalTransporter) {
|
||||
enqueEvent(listener -> listener.raceRemoved(trackedRace), threadLocalTransporter);
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -201,7 +215,7 @@ public class TrackedRegattaImpl implements TrackedRegatta {
|
||||
}
|
||||
}
|
||||
};
|
||||
addRaceListener(listener);
|
||||
addRaceListener(listener, Optional.empty());
|
||||
try {
|
||||
synchronized (mutex) {
|
||||
result = getExistingTrackedRace(race);
|
||||
@@ -232,16 +246,19 @@ public class TrackedRegattaImpl implements TrackedRegatta {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addRaceListener(RaceListener listener) {
|
||||
public void addRaceListener(RaceListener listener, Optional<ThreadLocalTransporter> threadLocalTransporter) {
|
||||
lockTrackedRacesForRead();
|
||||
try {
|
||||
raceListeners.put(listener, listener);
|
||||
final List<TrackedRace> trackedRacesCopy = new ArrayList<>();
|
||||
Util.addAll(getTrackedRaces(), trackedRacesCopy);
|
||||
threadLocalTransporter.ifPresent(ThreadLocalTransporter::rememberThreadLocalStates);
|
||||
eventQueue.addWork(() -> {
|
||||
for (TrackedRace trackedRace : trackedRacesCopy) {
|
||||
listener.raceAdded(trackedRace);
|
||||
}
|
||||
withBeforeAndAfterHandling(threadLocalTransporter, () -> {
|
||||
for (TrackedRace trackedRace : trackedRacesCopy) {
|
||||
listener.raceAdded(trackedRace);
|
||||
}
|
||||
});
|
||||
});
|
||||
} finally {
|
||||
unlockTrackedRacesAfterRead();
|
||||
@@ -284,7 +301,8 @@ public class TrackedRegattaImpl implements TrackedRegatta {
|
||||
public DynamicTrackedRace createTrackedRace(RaceDefinition raceDefinition, Iterable<Sideline> sidelines,
|
||||
WindStore windStore, long delayToLiveInMillis,
|
||||
long millisecondsOverWhichToAverageWind, long millisecondsOverWhichToAverageSpeed,
|
||||
DynamicRaceDefinitionSet raceDefinitionSetToUpdate, boolean useInternalMarkPassingAlgorithm, RaceLogResolver raceLogResolver) {
|
||||
DynamicRaceDefinitionSet raceDefinitionSetToUpdate, boolean useInternalMarkPassingAlgorithm, RaceLogResolver raceLogResolver,
|
||||
Optional<ThreadLocalTransporter> threadLocalTransporter) {
|
||||
logger.log(Level.INFO, "Creating DynamicTrackedRaceImpl for RaceDefinition " + raceDefinition.getName());
|
||||
DynamicTrackedRaceImpl result = new DynamicTrackedRaceImpl(this, raceDefinition, sidelines, windStore,
|
||||
delayToLiveInMillis, millisecondsOverWhichToAverageWind,
|
||||
@@ -295,7 +313,7 @@ public class TrackedRegattaImpl implements TrackedRegatta {
|
||||
if (raceDefinitionSetToUpdate != null) {
|
||||
raceDefinitionSetToUpdate.addRaceDefinition(raceDefinition, result);
|
||||
}
|
||||
addTrackedRace(result);
|
||||
addTrackedRace(result, threadLocalTransporter);
|
||||
return result;
|
||||
}
|
||||
}
|
||||
|
||||
+6
-4
@@ -7,6 +7,7 @@ import java.io.Serializable;
|
||||
import java.util.Collections;
|
||||
import java.util.List;
|
||||
import java.util.NavigableSet;
|
||||
import java.util.Optional;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.Future;
|
||||
|
||||
@@ -76,6 +77,7 @@ import com.sap.sse.common.Duration;
|
||||
import com.sap.sse.common.IsManagedByCache;
|
||||
import com.sap.sse.common.TimePoint;
|
||||
import com.sap.sse.common.Util;
|
||||
import com.sap.sse.util.ThreadLocalTransporter;
|
||||
|
||||
public class MockedTrackedRace implements DynamicTrackedRace {
|
||||
private static final long serialVersionUID = 5827912985564121181L;
|
||||
@@ -295,15 +297,15 @@ public class MockedTrackedRace implements DynamicTrackedRace {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addTrackedRace(TrackedRace trackedRace) {
|
||||
public void addTrackedRace(TrackedRace trackedRace, Optional<ThreadLocalTransporter> threadLocalTransporter) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void removeTrackedRace(TrackedRace trackedRace) {
|
||||
public void removeTrackedRace(TrackedRace trackedRace, Optional<ThreadLocalTransporter> threadLocalTransporter) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addRaceListener(RaceListener listener) {
|
||||
public void addRaceListener(RaceListener listener, Optional<ThreadLocalTransporter> threadLocalTransporter) {
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -325,7 +327,7 @@ public class MockedTrackedRace implements DynamicTrackedRace {
|
||||
public DynamicTrackedRace createTrackedRace(RaceDefinition raceDefinition, Iterable<Sideline> sidelines, WindStore windStore,
|
||||
long delayToLiveInMillis, long millisecondsOverWhichToAverageWind,
|
||||
long millisecondsOverWhichToAverageSpeed, DynamicRaceDefinitionSet raceDefinitionSetToUpdate,
|
||||
boolean useMarkPassingCalculator, RaceLogResolver raceLogResolver) {
|
||||
boolean useMarkPassingCalculator, RaceLogResolver raceLogResolver, Optional<ThreadLocalTransporter> threadLocalTransporter) {
|
||||
return null;
|
||||
}
|
||||
|
||||
|
||||
+3
-1
@@ -9,6 +9,7 @@ import java.net.MalformedURLException;
|
||||
import java.net.URISyntaxException;
|
||||
import java.net.UnknownHostException;
|
||||
import java.util.Collections;
|
||||
import java.util.Optional;
|
||||
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
@@ -103,7 +104,8 @@ public class TestStoringAndRetrievingWindTracksTest extends AbstractTracTracLive
|
||||
@Override
|
||||
public void addRaceDefinition(RaceDefinition race, DynamicTrackedRace trackedRace) {
|
||||
}
|
||||
}, /*useMarkPassingCalculator*/ false, mock(RaceLogResolver.class));
|
||||
}, /*useMarkPassingCalculator*/ false, mock(RaceLogResolver.class),
|
||||
Optional.empty());
|
||||
WindSource windSource = new WindSourceImpl(WindSourceType.WEB);
|
||||
Mongo myFirstMongo = newMongo();
|
||||
DB firstDatabase = myFirstMongo.getDB(dbConfiguration.getDatabaseName());
|
||||
|
||||
+3
-2
@@ -8,6 +8,7 @@ import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.Collections;
|
||||
import java.util.List;
|
||||
import java.util.Optional;
|
||||
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
@@ -81,7 +82,7 @@ public class StatisticsTest {
|
||||
trackedRace.setEndOfTrackingReceived(new MillisecondsTimePoint(END_OF_TRACKING));
|
||||
trackedRace.setStartTimeReceived(new MillisecondsTimePoint(START_OF_RACE));
|
||||
|
||||
regatta.addTrackedRace(trackedRace);
|
||||
regatta.addTrackedRace(trackedRace, Optional.empty());
|
||||
}
|
||||
|
||||
private TrackedRaceStatisticsCacheImpl getStatisticsCacheWithRegattaAdded() throws Exception {
|
||||
@@ -98,7 +99,7 @@ public class StatisticsTest {
|
||||
public void raceAdded(TrackedRace trackedRace) {
|
||||
}
|
||||
};
|
||||
regatta.addRaceListener(raceListener);
|
||||
regatta.addRaceListener(raceListener, Optional.empty());
|
||||
regatta.removeRaceListener(raceListener).get();
|
||||
|
||||
return trackedRaceStatisticsCache;
|
||||
|
||||
+7
-3
@@ -14,6 +14,7 @@ import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.Collections;
|
||||
import java.util.Iterator;
|
||||
import java.util.Optional;
|
||||
import java.util.Set;
|
||||
import java.util.UUID;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
@@ -83,17 +84,20 @@ public class RaceTrackerStartStopTest {
|
||||
trackedRegatta1.createTrackedRace(raceDef1, Collections.<Sideline> emptyList(),
|
||||
/* windStore */ EmptyWindStore.INSTANCE, /* delayToLiveInMillis */ 0l,
|
||||
/* millisecondsOverWhichToAverageWind */ 0l,
|
||||
/* millisecondsOverWhichToAverageSpeed */ 0l, /* raceDefinitionSetToUpdate */ null, /*useMarkPassingCalculator*/ false, mock(RaceLogResolver.class));
|
||||
/* millisecondsOverWhichToAverageSpeed */ 0l, /* raceDefinitionSetToUpdate */ null, /*useMarkPassingCalculator*/ false, mock(RaceLogResolver.class),
|
||||
Optional.empty());
|
||||
regatta.addRace(raceDef2);
|
||||
trackedRegatta1.createTrackedRace(raceDef2, Collections.<Sideline> emptyList(),
|
||||
/* windStore */ EmptyWindStore.INSTANCE, /* delayToLiveInMillis */ 0l,
|
||||
/* millisecondsOverWhichToAverageWind */ 0l,
|
||||
/* millisecondsOverWhichToAverageSpeed */ 0l, /* raceDefinitionSetToUpdate */ null, /*useMarkPassingCalculator*/ false, mock(RaceLogResolver.class));
|
||||
/* millisecondsOverWhichToAverageSpeed */ 0l, /* raceDefinitionSetToUpdate */ null, /*useMarkPassingCalculator*/ false, mock(RaceLogResolver.class),
|
||||
Optional.empty());
|
||||
regatta.addRace(raceDef3);
|
||||
trackedRegatta1.createTrackedRace(raceDef3, Collections.<Sideline> emptyList(),
|
||||
/* windStore */ EmptyWindStore.INSTANCE, /* delayToLiveInMillis */ 0l,
|
||||
/* millisecondsOverWhichToAverageWind */ 0l,
|
||||
/* millisecondsOverWhichToAverageSpeed */ 0l, /* raceDefinitionSetToUpdate */ null, /*useMarkPassingCalculator*/ false, mock(RaceLogResolver.class));
|
||||
/* millisecondsOverWhichToAverageSpeed */ 0l, /* raceDefinitionSetToUpdate */ null, /*useMarkPassingCalculator*/ false, mock(RaceLogResolver.class),
|
||||
Optional.empty());
|
||||
Long trackerID1 = new Long(1);
|
||||
Long trackerID2 = new Long(2);
|
||||
Long trackerID3 = new Long(3);
|
||||
|
||||
+2
-1
@@ -10,6 +10,7 @@ import java.net.MalformedURLException;
|
||||
import java.net.URI;
|
||||
import java.net.URISyntaxException;
|
||||
import java.net.URL;
|
||||
import java.util.Optional;
|
||||
import java.util.logging.Logger;
|
||||
|
||||
import org.junit.After;
|
||||
@@ -102,7 +103,7 @@ public class RaceTrackerTest {
|
||||
@Override
|
||||
public void raceRemoved(TrackedRace trackedRace) {
|
||||
}
|
||||
});
|
||||
}, Optional.empty());
|
||||
synchronized (trackedRaces) {
|
||||
if (trackedRaces[0] == null) {
|
||||
trackedRaces.wait();
|
||||
|
||||
+3
-1
@@ -10,6 +10,7 @@ import java.util.Collections;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.UUID;
|
||||
|
||||
import org.junit.Before;
|
||||
@@ -101,7 +102,8 @@ public class RemoveLeaderboardTest {
|
||||
trackedRace = trackedRegatta1.createTrackedRace(raceDef1, Collections.<Sideline> emptyList(),
|
||||
/* windStore */ EmptyWindStore.INSTANCE, /* delayToLiveInMillis */ 0l,
|
||||
/* millisecondsOverWhichToAverageWind */ 0l,
|
||||
/* millisecondsOverWhichToAverageSpeed */ 0l, /* raceDefinitionSetToUpdate */ null, /*useMarkPassingCalculator*/ false, mock(RaceLogResolver.class));
|
||||
/* millisecondsOverWhichToAverageSpeed */ 0l, /* raceDefinitionSetToUpdate */ null, /*useMarkPassingCalculator*/ false, mock(RaceLogResolver.class),
|
||||
Optional.empty());
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
+7
-3
@@ -25,6 +25,7 @@ import java.util.List;
|
||||
import java.util.Locale;
|
||||
import java.util.Map;
|
||||
import java.util.Map.Entry;
|
||||
import java.util.Optional;
|
||||
import java.util.Set;
|
||||
import java.util.UUID;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
@@ -1740,13 +1741,15 @@ public class RacingEventServiceImpl implements RacingEventService, ClearStateTes
|
||||
RaceDefinition race = getRace(raceIdentifier);
|
||||
return trackedRegatta.createTrackedRace(race, Collections.<Sideline> emptyList(), windStore,
|
||||
delayToLiveInMillis, millisecondsOverWhichToAverageWind, millisecondsOverWhichToAverageSpeed,
|
||||
/* raceDefinitionSetToUpdate */null, useMarkPassingCalculator, /* raceLogResolver */ this);
|
||||
/* raceDefinitionSetToUpdate */null, useMarkPassingCalculator, /* raceLogResolver */ this,
|
||||
Optional.of(this.getThreadLocalTransporterForCurrentlyFillingFromInitialLoadOrApplyingOperationReceivedFromMaster()));
|
||||
}
|
||||
|
||||
private void ensureRegattaIsObservedForDefaultLeaderboardAndAutoLeaderboardLinking(
|
||||
DynamicTrackedRegatta trackedRegatta) {
|
||||
if (regattasObservedForDefaultLeaderboard.add(trackedRegatta)) {
|
||||
trackedRegatta.addRaceListener(new RaceAdditionListener());
|
||||
trackedRegatta.addRaceListener(new RaceAdditionListener(),
|
||||
Optional.of(this.getThreadLocalTransporterForCurrentlyFillingFromInitialLoadOrApplyingOperationReceivedFromMaster()));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2436,7 +2439,8 @@ public class RacingEventServiceImpl implements RacingEventService, ClearStateTes
|
||||
final int newSizeOfTrackedRaces;
|
||||
oldSizeOfTrackedRaces = Util.size(trackedRegatta.getTrackedRaces());
|
||||
try {
|
||||
trackedRegatta.removeTrackedRace(trackedRace);
|
||||
trackedRegatta.removeTrackedRace(trackedRace, Optional.of(
|
||||
getThreadLocalTransporterForCurrentlyFillingFromInitialLoadOrApplyingOperationReceivedFromMaster()));
|
||||
newSizeOfTrackedRaces = Util.size(trackedRegatta.getTrackedRaces());
|
||||
isTrackedRacesBecameEmpty = (oldSizeOfTrackedRaces > 0 && newSizeOfTrackedRaces == 0);
|
||||
} finally {
|
||||
|
||||
+3
-1
@@ -7,6 +7,7 @@ import java.util.Collection;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Map.Entry;
|
||||
import java.util.Optional;
|
||||
import java.util.Set;
|
||||
import java.util.UUID;
|
||||
import java.util.logging.Level;
|
||||
@@ -477,7 +478,8 @@ public class ImportMasterDataOperation extends
|
||||
trackedRegatta.unlockTrackedRacesAfterRead();
|
||||
}
|
||||
for (TrackedRace raceToRemove : toRemove) {
|
||||
trackedRegatta.removeTrackedRace(raceToRemove);
|
||||
trackedRegatta.removeTrackedRace(raceToRemove, Optional.of(toState
|
||||
.getThreadLocalTransporterForCurrentlyFillingFromInitialLoadOrApplyingOperationReceivedFromMaster()));
|
||||
RaceDefinition race = existingRegatta.getRaceByName(raceToRemove
|
||||
.getRaceIdentifier().getRaceName());
|
||||
if (race != null) {
|
||||
|
||||
+2
-1
@@ -6,6 +6,7 @@ import java.util.Iterator;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Map.Entry;
|
||||
import java.util.Optional;
|
||||
import java.util.TreeMap;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import java.util.concurrent.Future;
|
||||
@@ -283,7 +284,7 @@ public class SimulationServiceImpl implements SimulationService {
|
||||
DynamicTrackedRegatta trackedRegatta = racingEventService.getTrackedRegatta(regatta);
|
||||
SimulationRaceListener raceListener = new SimulationRaceListener();
|
||||
raceListeners.put(legIdentifier.getRegattaName(), raceListener);
|
||||
trackedRegatta.addRaceListener(raceListener);
|
||||
trackedRegatta.addRaceListener(raceListener, /* Not replicated */ Optional.empty());
|
||||
}
|
||||
if (!legListeners.containsKey(legIdentifier.getRaceIdentifier())) {
|
||||
TrackedRace trackedRace = racingEventService.getTrackedRace(legIdentifier);
|
||||
|
||||
Reference in New Issue
Block a user